Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
5c01911e3c | ||
|
|
6041479286 |
+8
-97
@@ -4,7 +4,6 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/nats-io/nats.go"
|
"github.com/nats-io/nats.go"
|
||||||
)
|
)
|
||||||
@@ -138,43 +137,11 @@ func (dvm *DistributedVM) LoadModel(model RuntimeModel) error {
|
|||||||
|
|
||||||
// SendMessage routes messages across the distributed cluster
|
// SendMessage routes messages across the distributed cluster
|
||||||
func (dvm *DistributedVM) SendMessage(message RuntimeMessage) error {
|
func (dvm *DistributedVM) SendMessage(message RuntimeMessage) error {
|
||||||
actorID := message.GetTargetActorID()
|
// This is a simplified implementation
|
||||||
targetNode := dvm.GetActorNode(actorID)
|
// In practice, this would determine the target node based on sharding
|
||||||
|
// and route the message appropriately
|
||||||
|
|
||||||
if targetNode == dvm.nodeID {
|
|
||||||
return dvm.localRuntime.SendMessage(message)
|
return dvm.localRuntime.SendMessage(message)
|
||||||
}
|
|
||||||
|
|
||||||
return dvm.routeMessageToNode(actorID, message)
|
|
||||||
}
|
|
||||||
|
|
||||||
// routeMessageToNode sends a message to another node for delivery to the target actor
|
|
||||||
func (dvm *DistributedVM) routeMessageToNode(actorID string, message RuntimeMessage) error {
|
|
||||||
hops := 0
|
|
||||||
var body map[string]interface{}
|
|
||||||
if mp, ok := message.(*MessagePayload); ok {
|
|
||||||
hops = mp.Hops
|
|
||||||
body = mp.Body
|
|
||||||
}
|
|
||||||
if hops >= MaxRouteHops {
|
|
||||||
dvm.cluster.logger.Printf("Dropping message for actor %s: exceeded max hops (%d)", actorID, MaxRouteHops)
|
|
||||||
return fmt.Errorf("message exceeded max hops")
|
|
||||||
}
|
|
||||||
|
|
||||||
msg := ClusterMessage{
|
|
||||||
Type: "route_message",
|
|
||||||
From: dvm.nodeID,
|
|
||||||
To: actorID,
|
|
||||||
Payload: MessagePayload{
|
|
||||||
TargetActorID: actorID,
|
|
||||||
Type: message.GetType(),
|
|
||||||
Hops: hops + 1,
|
|
||||||
Body: body,
|
|
||||||
},
|
|
||||||
Timestamp: time.Now(),
|
|
||||||
}
|
|
||||||
|
|
||||||
return dvm.publishClusterMessage(msg)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// GetActorNode determines which node should handle a specific actor
|
// GetActorNode determines which node should handle a specific actor
|
||||||
@@ -222,10 +189,8 @@ func (dvm *DistributedVM) handleClusterMessage(msg *nats.Msg) {
|
|||||||
dvm.localRuntime.LoadModel(&model)
|
dvm.localRuntime.LoadModel(&model)
|
||||||
|
|
||||||
case "route_message":
|
case "route_message":
|
||||||
if clusterMsg.From == dvm.nodeID {
|
// Handle message routing from other nodes
|
||||||
return
|
// Re-marshal and unmarshal to convert map[string]interface{} to concrete type
|
||||||
}
|
|
||||||
|
|
||||||
payloadBytes, err := json.Marshal(clusterMsg.Payload)
|
payloadBytes, err := json.Marshal(clusterMsg.Payload)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return
|
return
|
||||||
@@ -234,24 +199,7 @@ func (dvm *DistributedVM) handleClusterMessage(msg *nats.Msg) {
|
|||||||
if err := json.Unmarshal(payloadBytes, &message); err != nil {
|
if err := json.Unmarshal(payloadBytes, &message); err != nil {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
dvm.localRuntime.SendMessage(&message)
|
||||||
if message.Hops >= MaxRouteHops {
|
|
||||||
dvm.cluster.logger.Printf("Dropping message for actor %s: exceeded max hops (%d)", message.TargetActorID, MaxRouteHops)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
targetActor := message.TargetActorID
|
|
||||||
msg := &MessagePayload{
|
|
||||||
TargetActorID: targetActor,
|
|
||||||
Type: message.Type,
|
|
||||||
Hops: message.Hops,
|
|
||||||
Body: message.Body,
|
|
||||||
}
|
|
||||||
if dvm.IsLocalActor(targetActor) {
|
|
||||||
dvm.localRuntime.SendMessage(msg)
|
|
||||||
} else {
|
|
||||||
dvm.routeMessageToNode(targetActor, msg)
|
|
||||||
}
|
|
||||||
|
|
||||||
case "rebalance":
|
case "rebalance":
|
||||||
// Handle shard rebalancing requests
|
// Handle shard rebalancing requests
|
||||||
@@ -261,45 +209,8 @@ func (dvm *DistributedVM) handleClusterMessage(msg *nats.Msg) {
|
|||||||
|
|
||||||
// handleRebalanceRequest processes shard rebalancing requests
|
// handleRebalanceRequest processes shard rebalancing requests
|
||||||
func (dvm *DistributedVM) handleRebalanceRequest(msg ClusterMessage) {
|
func (dvm *DistributedVM) handleRebalanceRequest(msg ClusterMessage) {
|
||||||
if msg.From == dvm.nodeID {
|
// Simplified rebalancing logic
|
||||||
return
|
// In practice, this would implement complex actor migration
|
||||||
}
|
|
||||||
|
|
||||||
if !dvm.cluster.IsLeader() {
|
|
||||||
dvm.cluster.logger.Printf("Ignoring rebalance request: not the leader")
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
if dvm.cluster.shardMap == nil {
|
|
||||||
dvm.cluster.logger.Printf("Shard map is nil, skipping rebalance")
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
payloadBytes, err := json.Marshal(msg.Payload)
|
|
||||||
if err != nil {
|
|
||||||
dvm.cluster.logger.Printf("Failed to marshal rebalance payload: %v", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
var newShardMap ShardMap
|
|
||||||
if err := json.Unmarshal(payloadBytes, &newShardMap); err != nil {
|
|
||||||
dvm.cluster.logger.Printf("Failed to unmarshal shard map: %v", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
dvm.cluster.mutex.Lock()
|
|
||||||
if newShardMap.Version > dvm.cluster.shardMap.Version {
|
|
||||||
dvm.cluster.shardMap = &newShardMap
|
|
||||||
dvm.cluster.hashRing = NewConsistentHashRing()
|
|
||||||
for nodeID := range newShardMap.Nodes {
|
|
||||||
dvm.cluster.hashRing.AddNode(nodeID)
|
|
||||||
}
|
|
||||||
dvm.cluster.logger.Printf("Applied new shard map (version %d) from rebalance", newShardMap.Version)
|
|
||||||
} else {
|
|
||||||
dvm.cluster.logger.Printf("Ignoring stale shard map (got version %d, current %d)",
|
|
||||||
newShardMap.Version, dvm.cluster.shardMap.Version)
|
|
||||||
}
|
|
||||||
dvm.cluster.mutex.Unlock()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// publishClusterMessage sends a message to other cluster nodes
|
// publishClusterMessage sends a message to other cluster nodes
|
||||||
|
|||||||
+12
-188
@@ -154,10 +154,6 @@ func (cm *ClusterManager) handleClusterMessage(msg *nats.Msg) {
|
|||||||
if update, ok := clusterMsg.Payload.(NodeUpdate); ok {
|
if update, ok := clusterMsg.Payload.(NodeUpdate); ok {
|
||||||
cm.handleNodeUpdate(update)
|
cm.handleNodeUpdate(update)
|
||||||
}
|
}
|
||||||
case "shard_map":
|
|
||||||
cm.handleShardMapUpdate(clusterMsg)
|
|
||||||
case "migration_update":
|
|
||||||
cm.handleMigrationUpdate(clusterMsg)
|
|
||||||
default:
|
default:
|
||||||
cm.logger.Printf("Unknown cluster message type: %s", clusterMsg.Type)
|
cm.logger.Printf("Unknown cluster message type: %s", clusterMsg.Type)
|
||||||
}
|
}
|
||||||
@@ -221,91 +217,16 @@ func (cm *ClusterManager) handleNodeUpdate(update NodeUpdate) {
|
|||||||
func (cm *ClusterManager) handleRebalanceRequest(msg ClusterMessage) {
|
func (cm *ClusterManager) handleRebalanceRequest(msg ClusterMessage) {
|
||||||
cm.logger.Printf("Handling rebalance request from %s", msg.From)
|
cm.logger.Printf("Handling rebalance request from %s", msg.From)
|
||||||
|
|
||||||
if !cm.IsLeader() {
|
// Implementation would handle the specific rebalancing logic
|
||||||
cm.logger.Printf("Ignoring rebalance request: not the leader")
|
// This is a simplified version
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
cm.mutex.RLock()
|
|
||||||
activeNodes := make(map[string]*NodeInfo)
|
|
||||||
for nodeID, nodeInfo := range cm.nodes {
|
|
||||||
if nodeInfo.Status == NodeStatusActive {
|
|
||||||
activeNodes[nodeID] = nodeInfo
|
|
||||||
}
|
|
||||||
}
|
|
||||||
cm.mutex.RUnlock()
|
|
||||||
|
|
||||||
if len(activeNodes) == 0 {
|
|
||||||
cm.logger.Printf("No active nodes for rebalancing")
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
placement := &ConsistentHashPlacement{}
|
|
||||||
newShardMap, err := placement.RebalanceShards(cm.shardMap, activeNodes)
|
|
||||||
if err != nil {
|
|
||||||
cm.logger.Printf("Failed to compute new shard map: %v", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
cm.mutex.Lock()
|
|
||||||
cm.shardMap = newShardMap
|
|
||||||
cm.mutex.Unlock()
|
|
||||||
|
|
||||||
cm.hashRing = NewConsistentHashRing()
|
|
||||||
for nodeID := range activeNodes {
|
|
||||||
cm.hashRing.AddNode(nodeID)
|
|
||||||
}
|
|
||||||
|
|
||||||
cm.broadcastShardMap(newShardMap)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// handleMigrationRequest processes actor migration requests
|
// handleMigrationRequest processes actor migration requests
|
||||||
func (cm *ClusterManager) handleMigrationRequest(msg ClusterMessage) {
|
func (cm *ClusterManager) handleMigrationRequest(msg ClusterMessage) {
|
||||||
cm.logger.Printf("Handling migration request from %s", msg.From)
|
cm.logger.Printf("Handling migration request from %s", msg.From)
|
||||||
|
|
||||||
var migration ActorMigration
|
// Implementation would handle the specific migration logic
|
||||||
payloadBytes, err := json.Marshal(msg.Payload)
|
// This is a simplified version
|
||||||
if err != nil {
|
|
||||||
cm.logger.Printf("Failed to marshal migration payload: %v", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
if err := json.Unmarshal(payloadBytes, &migration); err != nil {
|
|
||||||
cm.logger.Printf("Failed to unmarshal migration request: %v", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
cm.logger.Printf("Actor %s migrating from %s to %s (shard %d)",
|
|
||||||
migration.ActorID, migration.FromNode, migration.ToNode, migration.ShardID)
|
|
||||||
|
|
||||||
if migration.FromNode == cm.nodeID {
|
|
||||||
cm.logger.Printf("Initiating local actor state export for %s", migration.ActorID)
|
|
||||||
migration.Status = string(MigrationInProgress)
|
|
||||||
cm.broadcastMigrationUpdate(migration)
|
|
||||||
}
|
|
||||||
|
|
||||||
if migration.ToNode == cm.nodeID {
|
|
||||||
cm.logger.Printf("Actor %s assigned to this node, waiting for state import", migration.ActorID)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// broadcastMigrationUpdate propagates migration status updates to the cluster
|
|
||||||
func (cm *ClusterManager) broadcastMigrationUpdate(migration ActorMigration) {
|
|
||||||
msg := ClusterMessage{
|
|
||||||
Type: "migration_update",
|
|
||||||
From: cm.nodeID,
|
|
||||||
To: "broadcast",
|
|
||||||
Payload: migration,
|
|
||||||
}
|
|
||||||
|
|
||||||
data, err := json.Marshal(msg)
|
|
||||||
if err != nil {
|
|
||||||
cm.logger.Printf("Failed to marshal migration update: %v", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
if err := cm.natsConn.Publish("aether.cluster.migration_update", data); err != nil {
|
|
||||||
cm.logger.Printf("Failed to publish migration update: %v", err)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// triggerShardRebalancing initiates shard rebalancing across the cluster
|
// triggerShardRebalancing initiates shard rebalancing across the cluster
|
||||||
@@ -316,11 +237,12 @@ func (cm *ClusterManager) triggerShardRebalancing(reason string) {
|
|||||||
|
|
||||||
cm.logger.Printf("Triggering shard rebalancing: %s", reason)
|
cm.logger.Printf("Triggering shard rebalancing: %s", reason)
|
||||||
|
|
||||||
|
// Get active nodes
|
||||||
|
var activeNodes []*NodeInfo
|
||||||
cm.mutex.RLock()
|
cm.mutex.RLock()
|
||||||
activeNodes := make(map[string]*NodeInfo)
|
for _, node := range cm.nodes {
|
||||||
for nodeID, nodeInfo := range cm.nodes {
|
if node.Status == NodeStatusActive {
|
||||||
if nodeInfo.Status == NodeStatusActive {
|
activeNodes = append(activeNodes, node)
|
||||||
activeNodes[nodeID] = nodeInfo
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
cm.mutex.RUnlock()
|
cm.mutex.RUnlock()
|
||||||
@@ -330,23 +252,8 @@ func (cm *ClusterManager) triggerShardRebalancing(reason string) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
placement := &ConsistentHashPlacement{}
|
// This would implement the actual rebalancing logic
|
||||||
newShardMap, err := placement.RebalanceShards(cm.shardMap, activeNodes)
|
cm.logger.Printf("Would rebalance across %d active nodes", len(activeNodes))
|
||||||
if err != nil {
|
|
||||||
cm.logger.Printf("Failed to compute new shard map: %v", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
cm.mutex.Lock()
|
|
||||||
cm.shardMap = newShardMap
|
|
||||||
cm.mutex.Unlock()
|
|
||||||
|
|
||||||
cm.hashRing = NewConsistentHashRing()
|
|
||||||
for nodeID := range activeNodes {
|
|
||||||
cm.hashRing.AddNode(nodeID)
|
|
||||||
}
|
|
||||||
|
|
||||||
cm.broadcastShardMap(newShardMap)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// monitorNodes periodically checks node health and updates
|
// monitorNodes periodically checks node health and updates
|
||||||
@@ -412,99 +319,16 @@ func (cm *ClusterManager) GetNodes() map[string]*NodeInfo {
|
|||||||
return nodes
|
return nodes
|
||||||
}
|
}
|
||||||
|
|
||||||
// handleShardMapUpdate applies a new shard map received from the leader
|
|
||||||
func (cm *ClusterManager) handleShardMapUpdate(msg ClusterMessage) {
|
|
||||||
if msg.From == cm.nodeID {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
payloadBytes, err := json.Marshal(msg.Payload)
|
|
||||||
if err != nil {
|
|
||||||
cm.logger.Printf("Failed to marshal shard map payload: %v", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
var newShardMap ShardMap
|
|
||||||
if err := json.Unmarshal(payloadBytes, &newShardMap); err != nil {
|
|
||||||
cm.logger.Printf("Failed to unmarshal shard map: %v", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
cm.mutex.Lock()
|
|
||||||
if newShardMap.Version > cm.shardMap.Version {
|
|
||||||
cm.shardMap = &newShardMap
|
|
||||||
cm.hashRing = NewConsistentHashRing()
|
|
||||||
for nodeID := range newShardMap.Nodes {
|
|
||||||
cm.hashRing.AddNode(nodeID)
|
|
||||||
}
|
|
||||||
cm.logger.Printf("Applied new shard map (version %d)", newShardMap.Version)
|
|
||||||
} else {
|
|
||||||
cm.logger.Printf("Ignoring stale shard map (got version %d, current %d)",
|
|
||||||
newShardMap.Version, cm.shardMap.Version)
|
|
||||||
}
|
|
||||||
cm.mutex.Unlock()
|
|
||||||
}
|
|
||||||
|
|
||||||
// GetShardMap returns the current shard mapping
|
// GetShardMap returns the current shard mapping
|
||||||
func (cm *ClusterManager) GetShardMap() *ShardMap {
|
func (cm *ClusterManager) GetShardMap() *ShardMap {
|
||||||
cm.mutex.RLock()
|
cm.mutex.RLock()
|
||||||
defer cm.mutex.RUnlock()
|
defer cm.mutex.RUnlock()
|
||||||
|
|
||||||
// Return a copy to prevent external mutation
|
// Return a copy to prevent external mutation
|
||||||
copy := &ShardMap{
|
return &ShardMap{
|
||||||
Version: cm.shardMap.Version,
|
Version: cm.shardMap.Version,
|
||||||
Shards: make(map[int][]string),
|
Shards: make(map[int][]string),
|
||||||
Nodes: make(map[string]NodeInfo),
|
Nodes: make(map[string]NodeInfo),
|
||||||
UpdateTime: cm.shardMap.UpdateTime,
|
UpdateTime: cm.shardMap.UpdateTime,
|
||||||
}
|
}
|
||||||
|
|
||||||
for shardID, nodes := range cm.shardMap.Shards {
|
|
||||||
copy.Shards[shardID] = append([]string(nil), nodes...)
|
|
||||||
}
|
|
||||||
|
|
||||||
for nodeID, nodeInfo := range cm.shardMap.Nodes {
|
|
||||||
copy.Nodes[nodeID] = nodeInfo
|
|
||||||
}
|
|
||||||
|
|
||||||
return copy
|
|
||||||
}
|
|
||||||
|
|
||||||
// broadcastShardMap propagates a new shard map to all cluster nodes via NATS
|
|
||||||
func (cm *ClusterManager) broadcastShardMap(newShardMap *ShardMap) {
|
|
||||||
msg := ClusterMessage{
|
|
||||||
Type: "shard_map",
|
|
||||||
From: cm.nodeID,
|
|
||||||
To: "broadcast",
|
|
||||||
Payload: newShardMap,
|
|
||||||
Timestamp: time.Now(),
|
|
||||||
}
|
|
||||||
|
|
||||||
data, err := json.Marshal(msg)
|
|
||||||
if err != nil {
|
|
||||||
cm.logger.Printf("Failed to marshal shard map broadcast: %v", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
if err := cm.natsConn.Publish("aether.cluster.shard_map", data); err != nil {
|
|
||||||
cm.logger.Printf("Failed to publish shard map broadcast: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
cm.logger.Printf("Broadcast new shard map (version %d) to cluster", newShardMap.Version)
|
|
||||||
}
|
|
||||||
|
|
||||||
// handleMigrationUpdate processes migration status update messages from other nodes
|
|
||||||
func (cm *ClusterManager) handleMigrationUpdate(msg ClusterMessage) {
|
|
||||||
var migration ActorMigration
|
|
||||||
payloadBytes, err := json.Marshal(msg.Payload)
|
|
||||||
if err != nil {
|
|
||||||
cm.logger.Printf("Failed to marshal migration update payload: %v", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
if err := json.Unmarshal(payloadBytes, &migration); err != nil {
|
|
||||||
cm.logger.Printf("Failed to unmarshal migration update: %v", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
cm.logger.Printf("Migration update for actor %s: status=%s (from %s)",
|
|
||||||
migration.ActorID, migration.Status, msg.From)
|
|
||||||
}
|
}
|
||||||
|
|||||||
+18
-100
@@ -6,8 +6,6 @@ import (
|
|||||||
"fmt"
|
"fmt"
|
||||||
"hash"
|
"hash"
|
||||||
"hash/fnv"
|
"hash/fnv"
|
||||||
"sort"
|
|
||||||
"time"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// MigrationStatus tracks actor migration progress
|
// MigrationStatus tracks actor migration progress
|
||||||
@@ -182,116 +180,36 @@ func (sm *ShardManager) GetReplicationFactor() int {
|
|||||||
// ConsistentHashPlacement implements PlacementStrategy using consistent hashing
|
// ConsistentHashPlacement implements PlacementStrategy using consistent hashing
|
||||||
type ConsistentHashPlacement struct{}
|
type ConsistentHashPlacement struct{}
|
||||||
|
|
||||||
// PlaceActor places an actor using the consistent hash ring
|
// PlaceActor places an actor using consistent hashing
|
||||||
func (chp *ConsistentHashPlacement) PlaceActor(actorID string, shardMap *ShardMap, nodes map[string]*NodeInfo) (string, error) {
|
func (chp *ConsistentHashPlacement) PlaceActor(actorID string, shardMap *ShardMap, nodes map[string]*NodeInfo) (string, error) {
|
||||||
if len(nodes) == 0 {
|
if len(nodes) == 0 {
|
||||||
return "", fmt.Errorf("no nodes available for placement")
|
return "", fmt.Errorf("no nodes available for placement")
|
||||||
}
|
}
|
||||||
|
|
||||||
ring := NewConsistentHashRing()
|
// Simple consistent hash placement - in a real implementation,
|
||||||
|
// this would use the consistent hash ring
|
||||||
|
h := sha256.Sum256([]byte(actorID))
|
||||||
|
nodeIndex := binary.BigEndian.Uint32(h[:4]) % uint32(len(nodes))
|
||||||
|
|
||||||
|
i := 0
|
||||||
for nodeID := range nodes {
|
for nodeID := range nodes {
|
||||||
ring.AddNode(nodeID)
|
if i == int(nodeIndex) {
|
||||||
|
return nodeID, nil
|
||||||
|
}
|
||||||
|
i++
|
||||||
}
|
}
|
||||||
|
|
||||||
node := ring.GetNode(actorID)
|
// Fallback to first node
|
||||||
if node == "" {
|
|
||||||
sortedNodeIDs := make([]string, 0, len(nodes))
|
|
||||||
for nodeID := range nodes {
|
for nodeID := range nodes {
|
||||||
sortedNodeIDs = append(sortedNodeIDs, nodeID)
|
return nodeID, nil
|
||||||
}
|
|
||||||
sort.Strings(sortedNodeIDs)
|
|
||||||
return sortedNodeIDs[0], nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
return node, nil
|
return "", fmt.Errorf("failed to place actor")
|
||||||
}
|
}
|
||||||
|
|
||||||
// RebalanceShards redistributes shards across nodes using consistent hashing
|
// RebalanceShards rebalances shards across nodes
|
||||||
func (chp *ConsistentHashPlacement) RebalanceShards(currentMap *ShardMap, nodes map[string]*NodeInfo) (*ShardMap, error) {
|
func (chp *ConsistentHashPlacement) RebalanceShards(currentMap *ShardMap, nodes map[string]*NodeInfo) (*ShardMap, error) {
|
||||||
if len(nodes) == 0 {
|
// This is a simplified implementation
|
||||||
return nil, fmt.Errorf("no nodes available for rebalancing")
|
// In practice, this would implement sophisticated rebalancing logic
|
||||||
}
|
return currentMap, nil
|
||||||
|
|
||||||
ring := NewConsistentHashRing()
|
|
||||||
for nodeID := range nodes {
|
|
||||||
ring.AddNode(nodeID)
|
|
||||||
}
|
|
||||||
|
|
||||||
replicaCount := chp.deriveReplicaCount(currentMap)
|
|
||||||
|
|
||||||
newMap := &ShardMap{
|
|
||||||
Version: currentMap.Version + 1,
|
|
||||||
Shards: make(map[int][]string),
|
|
||||||
Nodes: make(map[string]NodeInfo),
|
|
||||||
UpdateTime: time.Now(),
|
|
||||||
}
|
|
||||||
|
|
||||||
for nodeID, nodeInfo := range nodes {
|
|
||||||
newMap.Nodes[nodeID] = *nodeInfo
|
|
||||||
}
|
|
||||||
|
|
||||||
for shardID := range currentMap.Shards {
|
|
||||||
primaryNode := ring.GetNode(fmt.Sprintf("shard-%d", shardID))
|
|
||||||
if primaryNode == "" {
|
|
||||||
sortedNodeIDs := make([]string, 0, len(nodes))
|
|
||||||
for nodeID := range nodes {
|
|
||||||
sortedNodeIDs = append(sortedNodeIDs, nodeID)
|
|
||||||
}
|
|
||||||
sort.Strings(sortedNodeIDs)
|
|
||||||
primaryNode = sortedNodeIDs[0]
|
|
||||||
}
|
|
||||||
|
|
||||||
var replicaNodes []string
|
|
||||||
candidates := make([]string, 0, len(nodes))
|
|
||||||
for nodeID := range nodes {
|
|
||||||
if nodeID != primaryNode {
|
|
||||||
candidates = append(candidates, nodeID)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
sort.Strings(candidates)
|
|
||||||
|
|
||||||
for i := 0; i < replicaCount && len(replicaNodes) < replicaCount; i++ {
|
|
||||||
node := ring.GetNode(fmt.Sprintf("shard-%d-replica-%d", shardID, i))
|
|
||||||
if node != "" && node != primaryNode {
|
|
||||||
found := false
|
|
||||||
for _, existing := range replicaNodes {
|
|
||||||
if existing == node {
|
|
||||||
found = true
|
|
||||||
break
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if !found {
|
|
||||||
replicaNodes = append(replicaNodes, node)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if len(replicaNodes) == 0 && len(candidates) > 0 {
|
|
||||||
replicaNodes = append(replicaNodes, candidates[0])
|
|
||||||
}
|
|
||||||
|
|
||||||
if len(replicaNodes) > replicaCount {
|
|
||||||
replicaNodes = replicaNodes[:replicaCount]
|
|
||||||
}
|
|
||||||
|
|
||||||
shardNodes := []string{primaryNode}
|
|
||||||
shardNodes = append(shardNodes, replicaNodes...)
|
|
||||||
newMap.Shards[shardID] = shardNodes
|
|
||||||
}
|
|
||||||
|
|
||||||
return newMap, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// deriveReplicaCount extracts the replication factor from the current shard map
|
|
||||||
func (chp *ConsistentHashPlacement) deriveReplicaCount(currentMap *ShardMap) int {
|
|
||||||
maxNodes := 0
|
|
||||||
for _, nodes := range currentMap.Shards {
|
|
||||||
if len(nodes) > maxNodes {
|
|
||||||
maxNodes = len(nodes)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if maxNodes <= 1 {
|
|
||||||
return 1
|
|
||||||
}
|
|
||||||
return maxNodes - 1
|
|
||||||
}
|
}
|
||||||
|
|||||||
+4
-20
@@ -650,8 +650,7 @@ func TestConsistentHashPlacement_RebalanceShards(t *testing.T) {
|
|||||||
placement := &ConsistentHashPlacement{}
|
placement := &ConsistentHashPlacement{}
|
||||||
currentMap := &ShardMap{
|
currentMap := &ShardMap{
|
||||||
Version: 1,
|
Version: 1,
|
||||||
Shards: map[int][]string{0: {"node-1"}, 1: {"node-1"}, 2: {"node-2"}},
|
Shards: map[int][]string{0: {"node-1"}},
|
||||||
Nodes: map[string]NodeInfo{},
|
|
||||||
}
|
}
|
||||||
nodes := map[string]*NodeInfo{
|
nodes := map[string]*NodeInfo{
|
||||||
"node-1": {ID: "node-1"},
|
"node-1": {ID: "node-1"},
|
||||||
@@ -663,24 +662,9 @@ func TestConsistentHashPlacement_RebalanceShards(t *testing.T) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
t.Errorf("unexpected error: %v", err)
|
t.Errorf("unexpected error: %v", err)
|
||||||
}
|
}
|
||||||
if result == nil {
|
// Current implementation returns unchanged map
|
||||||
t.Fatal("rebalance returned nil")
|
if result != currentMap {
|
||||||
}
|
t.Error("expected same map returned (simplified implementation)")
|
||||||
if result.Version != currentMap.Version+1 {
|
|
||||||
t.Errorf("expected version %d, got %d", currentMap.Version+1, result.Version)
|
|
||||||
}
|
|
||||||
if len(result.Shards) != len(currentMap.Shards) {
|
|
||||||
t.Errorf("expected %d shards, got %d", len(currentMap.Shards), len(result.Shards))
|
|
||||||
}
|
|
||||||
for shardID, shardNodes := range result.Shards {
|
|
||||||
if len(shardNodes) == 0 {
|
|
||||||
t.Errorf("shard %d has no nodes assigned", shardID)
|
|
||||||
}
|
|
||||||
for _, node := range shardNodes {
|
|
||||||
if _, exists := nodes[node]; !exists {
|
|
||||||
t.Errorf("shard %d assigned to unknown node %s", shardID, node)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -191,16 +191,11 @@ func (m *ModelPayload) GetID() string { return m.ID }
|
|||||||
// GetName implements RuntimeModel
|
// GetName implements RuntimeModel
|
||||||
func (m *ModelPayload) GetName() string { return m.Name }
|
func (m *ModelPayload) GetName() string { return m.Name }
|
||||||
|
|
||||||
// MaxRouteHops is the maximum number of hops a routed message can take before being dropped
|
|
||||||
const MaxRouteHops = 10
|
|
||||||
|
|
||||||
// MessagePayload is a concrete type for JSON-unmarshaling RuntimeMessage payloads.
|
// MessagePayload is a concrete type for JSON-unmarshaling RuntimeMessage payloads.
|
||||||
// Use this when receiving message data over the network.
|
// Use this when receiving message data over the network.
|
||||||
type MessagePayload struct {
|
type MessagePayload struct {
|
||||||
TargetActorID string `json:"targetActorId"`
|
TargetActorID string `json:"targetActorId"`
|
||||||
Type string `json:"type"`
|
Type string `json:"type"`
|
||||||
Hops int `json:"hops,omitempty"`
|
|
||||||
Body map[string]interface{} `json:"body,omitempty"`
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// GetTargetActorID implements RuntimeMessage
|
// GetTargetActorID implements RuntimeMessage
|
||||||
|
|||||||
@@ -0,0 +1,317 @@
|
|||||||
|
//go:build integration
|
||||||
|
// +build integration
|
||||||
|
|
||||||
|
package examples
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"log"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"git.flowmade.one/flowmade-one/aether"
|
||||||
|
"git.flowmade.one/flowmade-one/aether/store"
|
||||||
|
"github.com/nats-io/nats.go"
|
||||||
|
)
|
||||||
|
|
||||||
|
// CrossNodeBroadcasting demonstrates how to implement cross-node event broadcasting
|
||||||
|
// using NATSEventBus with JetStreamEventStore. This example shows how events persist
|
||||||
|
// to JetStream and are then broadcast to other nodes in the cluster via NATS.
|
||||||
|
//
|
||||||
|
// Key Concepts:
|
||||||
|
// 1. NATSEventBus wraps EventBus to add NATS publishing
|
||||||
|
// 2. JetStreamEventStore with broadcaster publishes EventStored to NATS
|
||||||
|
// 3. Other nodes receive these events via NATS subscription
|
||||||
|
// 4. Version cache is updated on remote events to maintain consistency
|
||||||
|
//
|
||||||
|
// Usage:
|
||||||
|
// go run examples/cross_node_broadcasting.go
|
||||||
|
func CrossNodeBroadcastingExample() {
|
||||||
|
// Connect to NATS
|
||||||
|
nc, err := nats.Connect(nats.DefaultURL)
|
||||||
|
if err != nil {
|
||||||
|
log.Fatalf("Failed to connect to NATS: %v", err)
|
||||||
|
}
|
||||||
|
defer nc.Close()
|
||||||
|
|
||||||
|
// Create NATS event bus (this will broadcast to all nodes)
|
||||||
|
natsBus, err := aether.NewNATSEventBus(nc)
|
||||||
|
if err != nil {
|
||||||
|
log.Fatalf("Failed to create NATS event bus: %v", err)
|
||||||
|
}
|
||||||
|
defer natsBus.Stop()
|
||||||
|
|
||||||
|
// Create JetStream event store WITH broadcaster
|
||||||
|
// This enables EventStored events to be published to NATS
|
||||||
|
store, err := store.NewJetStreamEventStoreWithBroadcaster(
|
||||||
|
nc,
|
||||||
|
"events",
|
||||||
|
natsBus,
|
||||||
|
"tenant-abc", // Optional namespace for isolation
|
||||||
|
)
|
||||||
|
if err != nil {
|
||||||
|
log.Fatalf("Failed to create event store: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Subscribe to EventStored events to update version cache
|
||||||
|
// This keeps the version cache synchronized across nodes
|
||||||
|
eventStoredCh := natsBus.SubscribeWithFilter(
|
||||||
|
"tenant-abc",
|
||||||
|
&aether.SubscriptionFilter{
|
||||||
|
EventTypes: []string{aether.EventTypeEventStored},
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
go func() {
|
||||||
|
for event := range eventStoredCh {
|
||||||
|
actorID := event.Data["actorId"].(string)
|
||||||
|
version := int64(event.Data["version"].(float64))
|
||||||
|
store.UpdateVersionCache(actorID, version)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
// Now save an event - it will be:
|
||||||
|
// 1. Persisted to JetStream
|
||||||
|
// 2. Published to NATS as EventStored
|
||||||
|
// 3. Received by other nodes via NATS
|
||||||
|
// 4. Used to update version cache
|
||||||
|
|
||||||
|
event := &aether.Event{
|
||||||
|
ID: "event-1",
|
||||||
|
EventType: "OrderPlaced",
|
||||||
|
ActorID: "order-123",
|
||||||
|
Version: 1,
|
||||||
|
Data: map[string]interface{}{
|
||||||
|
"total": 100.00,
|
||||||
|
"item": "widget",
|
||||||
|
},
|
||||||
|
Timestamp: time.Now(),
|
||||||
|
}
|
||||||
|
|
||||||
|
err = store.SaveEvent(event)
|
||||||
|
if err != nil {
|
||||||
|
log.Fatalf("Failed to save event: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
log.Println("Event saved to JetStream and broadcast to NATS")
|
||||||
|
|
||||||
|
// Other nodes in the cluster will receive this event via NATS
|
||||||
|
// and update their version cache accordingly
|
||||||
|
|
||||||
|
// Subscribe to events in this namespace
|
||||||
|
eventCh := natsBus.Subscribe("tenant-abc")
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
select {
|
||||||
|
case receivedEvent := <-eventCh:
|
||||||
|
log.Printf("Received event via NATS: %s (version %d)",
|
||||||
|
receivedEvent.EventType, receivedEvent.Version)
|
||||||
|
case <-ctx.Done():
|
||||||
|
log.Println("Timeout waiting for event")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// CrossNodeBroadcastingMultiNode demonstrates a multi-node cluster setup.
|
||||||
|
// This simulates multiple nodes connecting to the same NATS cluster.
|
||||||
|
// In production, these would be separate processes/machines.
|
||||||
|
//
|
||||||
|
// This example requires a running NATS server with JetStream enabled.
|
||||||
|
func CrossNodeBroadcastingMultiNode() {
|
||||||
|
// Connect to NATS
|
||||||
|
nc, err := nats.Connect(nats.DefaultURL)
|
||||||
|
if err != nil {
|
||||||
|
log.Fatalf("Failed to connect to NATS: %v", err)
|
||||||
|
}
|
||||||
|
defer nc.Close()
|
||||||
|
|
||||||
|
nodeID := "node-1"
|
||||||
|
log.Printf("Starting %s", nodeID)
|
||||||
|
|
||||||
|
// Each node creates its own NATS event bus and event store
|
||||||
|
natsBus, err := aether.NewNATSEventBus(nc)
|
||||||
|
if err != nil {
|
||||||
|
log.Fatalf("Failed to create NATS event bus: %v", err)
|
||||||
|
}
|
||||||
|
defer natsBus.Stop()
|
||||||
|
|
||||||
|
store, err := store.NewJetStreamEventStoreWithBroadcaster(
|
||||||
|
nc,
|
||||||
|
"events",
|
||||||
|
natsBus,
|
||||||
|
"tenant-abc",
|
||||||
|
)
|
||||||
|
if err != nil {
|
||||||
|
log.Fatalf("Failed to create event store: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Setup EventStored subscription for cache synchronization
|
||||||
|
eventStoredCh := natsBus.SubscribeWithFilter(
|
||||||
|
"tenant-abc",
|
||||||
|
&aether.SubscriptionFilter{
|
||||||
|
EventTypes: []string{aether.EventTypeEventStored},
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
go func() {
|
||||||
|
for event := range eventStoredCh {
|
||||||
|
actorID := event.Data["actorId"].(string)
|
||||||
|
version := int64(event.Data["version"].(float64))
|
||||||
|
store.UpdateVersionCache(actorID, version)
|
||||||
|
log.Printf("[%s] Received EventStored for %s v%d", nodeID, actorID, version)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
// Subscribe to actual events
|
||||||
|
eventCh := natsBus.Subscribe("tenant-abc")
|
||||||
|
log.Printf("[%s] Subscribed to tenant-abc", nodeID)
|
||||||
|
|
||||||
|
// Save an event
|
||||||
|
savedEvent := &aether.Event{
|
||||||
|
ID: fmt.Sprintf("event-%s-1", nodeID),
|
||||||
|
EventType: "OrderPlaced",
|
||||||
|
ActorID: "order-123",
|
||||||
|
Version: 1,
|
||||||
|
Data: map[string]interface{}{
|
||||||
|
"node": nodeID,
|
||||||
|
"total": 100.00,
|
||||||
|
},
|
||||||
|
Timestamp: time.Now(),
|
||||||
|
}
|
||||||
|
|
||||||
|
err = store.SaveEvent(savedEvent)
|
||||||
|
if err != nil {
|
||||||
|
log.Fatalf("[%s] Failed to save event: %v", nodeID, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
log.Printf("[%s] Saved event to JetStream", nodeID)
|
||||||
|
|
||||||
|
// Wait to receive the event (either from local or remote)
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
select {
|
||||||
|
case receivedEvent := <-eventCh:
|
||||||
|
log.Printf("[%s] Received: %s v%d from %s",
|
||||||
|
nodeID,
|
||||||
|
receivedEvent.EventType,
|
||||||
|
receivedEvent.Version,
|
||||||
|
receivedEvent.GetCorrelationID(),
|
||||||
|
)
|
||||||
|
case <-ctx.Done():
|
||||||
|
log.Printf("[%s] Timeout waiting for event", nodeID)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// DistributedOrderProcessing demonstrates a realistic scenario where
|
||||||
|
// multiple nodes process orders for the same actor. This shows how
|
||||||
|
// cross-node broadcasting ensures consistency.
|
||||||
|
func DistributedOrderProcessing() {
|
||||||
|
// Connect to NATS
|
||||||
|
nc, err := nats.Connect(nats.DefaultURL)
|
||||||
|
if err != nil {
|
||||||
|
log.Fatalf("Failed to connect to NATS: %v", err)
|
||||||
|
}
|
||||||
|
defer nc.Close()
|
||||||
|
|
||||||
|
// Node 1: Order creation
|
||||||
|
node1Bus, _ := aether.NewNATSEventBus(nc)
|
||||||
|
node1Store, _ := store.NewJetStreamEventStoreWithBroadcaster(
|
||||||
|
nc, "events", node1Bus, "orders",
|
||||||
|
)
|
||||||
|
defer node1Bus.Stop()
|
||||||
|
|
||||||
|
// Setup EventStored subscription for Node 1
|
||||||
|
eventStoredCh1 := node1Bus.SubscribeWithFilter(
|
||||||
|
"orders",
|
||||||
|
&aether.SubscriptionFilter{
|
||||||
|
EventTypes: []string{aether.EventTypeEventStored},
|
||||||
|
},
|
||||||
|
)
|
||||||
|
go func() {
|
||||||
|
for event := range eventStoredCh1 {
|
||||||
|
actorID := event.Data["actorId"].(string)
|
||||||
|
version := int64(event.Data["version"].(float64))
|
||||||
|
node1Store.UpdateVersionCache(actorID, version)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
// Node 2: Order processing (different node)
|
||||||
|
node2Bus, _ := aether.NewNATSEventBus(nc)
|
||||||
|
node2Store, _ := store.NewJetStreamEventStoreWithBroadcaster(
|
||||||
|
nc, "events", node2Bus, "orders",
|
||||||
|
)
|
||||||
|
defer node2Bus.Stop()
|
||||||
|
|
||||||
|
// Setup EventStored subscription for Node 2
|
||||||
|
eventStoredCh2 := node2Bus.SubscribeWithFilter(
|
||||||
|
"orders",
|
||||||
|
&aether.SubscriptionFilter{
|
||||||
|
EventTypes: []string{aether.EventTypeEventStored},
|
||||||
|
},
|
||||||
|
)
|
||||||
|
go func() {
|
||||||
|
for event := range eventStoredCh2 {
|
||||||
|
actorID := event.Data["actorId"].(string)
|
||||||
|
version := int64(event.Data["version"].(float64))
|
||||||
|
node2Store.UpdateVersionCache(actorID, version)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
// Node 2 also subscribes to events to receive updates
|
||||||
|
node2EventCh := node2Bus.Subscribe("orders")
|
||||||
|
|
||||||
|
// Node 1 creates an order
|
||||||
|
orderPlaced := &aether.Event{
|
||||||
|
ID: "order-created",
|
||||||
|
EventType: "OrderPlaced",
|
||||||
|
ActorID: "order-456",
|
||||||
|
Version: 1,
|
||||||
|
Data: map[string]interface{}{
|
||||||
|
"total": 100.00,
|
||||||
|
},
|
||||||
|
Timestamp: time.Now(),
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := node1Store.SaveEvent(orderPlaced); err != nil {
|
||||||
|
log.Fatalf("Failed to create order: %v", err)
|
||||||
|
}
|
||||||
|
log.Println("Node 1: Created order-456")
|
||||||
|
|
||||||
|
// Node 2 receives the OrderPlaced event via NATS
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
select {
|
||||||
|
case event := <-node2EventCh:
|
||||||
|
log.Printf("Node 2: Received %s v%d", event.EventType, event.Version)
|
||||||
|
|
||||||
|
// Node 2 processes the order (must use version 2)
|
||||||
|
orderProcessed := &aether.Event{
|
||||||
|
ID: "order-processed",
|
||||||
|
EventType: "OrderProcessed",
|
||||||
|
ActorID: "order-456",
|
||||||
|
Version: 2, // Must be > 1
|
||||||
|
Data: map[string]interface{}{
|
||||||
|
"status": "shipped",
|
||||||
|
},
|
||||||
|
Timestamp: time.Now(),
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := node2Store.SaveEvent(orderProcessed); err != nil {
|
||||||
|
log.Fatalf("Node 2: Failed to process order: %v", err)
|
||||||
|
}
|
||||||
|
log.Println("Node 2: Processed order-456")
|
||||||
|
|
||||||
|
case <-ctx.Done():
|
||||||
|
log.Println("Node 2: Timeout waiting for order event")
|
||||||
|
}
|
||||||
|
|
||||||
|
// Verify event stream consistency
|
||||||
|
events, err := node1Store.GetEvents("order-456", 0)
|
||||||
|
if err != nil {
|
||||||
|
log.Fatalf("Failed to get events: %v", err)
|
||||||
|
}
|
||||||
|
log.Printf("Node 1: Event stream has %d events", len(events))
|
||||||
|
}
|
||||||
@@ -6,16 +6,19 @@ require (
|
|||||||
github.com/google/uuid v1.6.0
|
github.com/google/uuid v1.6.0
|
||||||
github.com/nats-io/nats.go v1.37.0
|
github.com/nats-io/nats.go v1.37.0
|
||||||
github.com/prometheus/client_golang v1.23.2
|
github.com/prometheus/client_golang v1.23.2
|
||||||
|
github.com/stretchr/testify v1.11.1
|
||||||
)
|
)
|
||||||
|
|
||||||
require (
|
require (
|
||||||
github.com/beorn7/perks v1.0.1 // indirect
|
github.com/beorn7/perks v1.0.1 // indirect
|
||||||
github.com/cespare/xxhash/v2 v2.3.0 // indirect
|
github.com/cespare/xxhash/v2 v2.3.0 // indirect
|
||||||
|
github.com/davecgh/go-spew v1.1.1 // indirect
|
||||||
github.com/klauspost/compress v1.18.0 // indirect
|
github.com/klauspost/compress v1.18.0 // indirect
|
||||||
github.com/kr/text v0.2.0 // indirect
|
github.com/kr/text v0.2.0 // indirect
|
||||||
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
|
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
|
||||||
github.com/nats-io/nkeys v0.4.7 // indirect
|
github.com/nats-io/nkeys v0.4.7 // indirect
|
||||||
github.com/nats-io/nuid v1.0.1 // indirect
|
github.com/nats-io/nuid v1.0.1 // indirect
|
||||||
|
github.com/pmezard/go-difflib v1.0.0 // indirect
|
||||||
github.com/prometheus/client_model v0.6.2 // indirect
|
github.com/prometheus/client_model v0.6.2 // indirect
|
||||||
github.com/prometheus/common v0.66.1 // indirect
|
github.com/prometheus/common v0.66.1 // indirect
|
||||||
github.com/prometheus/procfs v0.16.1 // indirect
|
github.com/prometheus/procfs v0.16.1 // indirect
|
||||||
@@ -23,4 +26,5 @@ require (
|
|||||||
golang.org/x/crypto v0.18.0 // indirect
|
golang.org/x/crypto v0.18.0 // indirect
|
||||||
golang.org/x/sys v0.35.0 // indirect
|
golang.org/x/sys v0.35.0 // indirect
|
||||||
google.golang.org/protobuf v1.36.8 // indirect
|
google.golang.org/protobuf v1.36.8 // indirect
|
||||||
|
gopkg.in/yaml.v3 v3.0.1 // indirect
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -210,6 +210,26 @@ func (neb *NATSEventBus) Publish(namespaceID string, event *Event) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// SubscribeToEventStored creates a subscription specifically for EventStored events.
|
||||||
|
// This is a convenience method for the common pattern of listening to persisted events
|
||||||
|
// to update version cache or trigger other actions.
|
||||||
|
//
|
||||||
|
// Example:
|
||||||
|
//
|
||||||
|
// eventStoredCh := natsBus.SubscribeToEventStored("tenant-abc")
|
||||||
|
// go func() {
|
||||||
|
// for event := range eventStoredCh {
|
||||||
|
// actorID := event.Data["actorId"].(string)
|
||||||
|
// version := int64(event.Data["version"].(float64))
|
||||||
|
// store.UpdateVersionCache(actorID, version)
|
||||||
|
// }
|
||||||
|
// }()
|
||||||
|
func (neb *NATSEventBus) SubscribeToEventStored(namespacePattern string) <-chan *Event {
|
||||||
|
return neb.SubscribeWithFilter(namespacePattern, &SubscriptionFilter{
|
||||||
|
EventTypes: []string{EventTypeEventStored},
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
// Stop closes the NATS event bus and all subscriptions
|
// Stop closes the NATS event bus and all subscriptions
|
||||||
func (neb *NATSEventBus) Stop() {
|
func (neb *NATSEventBus) Stop() {
|
||||||
neb.mutex.Lock()
|
neb.mutex.Lock()
|
||||||
|
|||||||
@@ -0,0 +1,6 @@
|
|||||||
|
{
|
||||||
|
"$schema": "https://docs.renovatebot.com/renovate-schema.json",
|
||||||
|
"extends": [
|
||||||
|
"config:recommended"
|
||||||
|
]
|
||||||
|
}
|
||||||
@@ -0,0 +1,36 @@
|
|||||||
|
//go:build integration
|
||||||
|
// +build integration
|
||||||
|
|
||||||
|
package store
|
||||||
|
|
||||||
|
import (
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/nats-io/nats.go"
|
||||||
|
)
|
||||||
|
|
||||||
|
// getTestNATSConnection returns a NATS connection for testing.
|
||||||
|
// This helper is used by benchmark tests that require NATS.
|
||||||
|
func getTestNATSConnection(t *testing.T) *nats.Conn {
|
||||||
|
nc, err := nats.Connect(nats.DefaultURL)
|
||||||
|
if err != nil {
|
||||||
|
t.Skipf("NATS not available: %v", err)
|
||||||
|
}
|
||||||
|
return nc
|
||||||
|
}
|
||||||
|
|
||||||
|
// getVersionFromEvent extracts version from EventStored event data.
|
||||||
|
func getVersionFromEvent(data map[string]interface{}) int64 {
|
||||||
|
switch v := data["version"].(type) {
|
||||||
|
case float64:
|
||||||
|
return int64(v)
|
||||||
|
case int64:
|
||||||
|
return v
|
||||||
|
case int:
|
||||||
|
return int64(v)
|
||||||
|
case uint64:
|
||||||
|
return int64(v)
|
||||||
|
default:
|
||||||
|
return 0
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,420 @@
|
|||||||
|
//go:build integration
|
||||||
|
// +build integration
|
||||||
|
|
||||||
|
package store
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"log"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"git.flowmade.one/flowmade-one/aether"
|
||||||
|
"github.com/nats-io/nats.go"
|
||||||
|
"github.com/stretchr/testify/assert"
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
)
|
||||||
|
|
||||||
|
// generateStreamName creates a unique stream name for each test run
|
||||||
|
func generateStreamName(baseName string) string {
|
||||||
|
return fmt.Sprintf("%s_%s_%d", baseName, "tv149", time.Now().UnixNano()%100000000)
|
||||||
|
}
|
||||||
|
|
||||||
|
// cleanupStream deletes a JetStream stream if it exists.
|
||||||
|
func cleanupStream(nc *nats.Conn, streamName string) {
|
||||||
|
js, err := nc.JetStream()
|
||||||
|
if err != nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
err = js.DeleteStream(streamName)
|
||||||
|
// Silently ignore errors - we just want to clean up
|
||||||
|
_ = err
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestCrossNodeBroadcasting_SingleNode tests basic cross-node broadcasting
|
||||||
|
// on a single node (local loopback).
|
||||||
|
func TestCrossNodeBroadcasting_SingleNode(t *testing.T) {
|
||||||
|
nc, err := nats.Connect(nats.DefaultURL)
|
||||||
|
require.NoError(t, err)
|
||||||
|
defer nc.Close()
|
||||||
|
|
||||||
|
streamName := generateStreamName("broadcast_single")
|
||||||
|
cleanupStream(nc, streamName)
|
||||||
|
|
||||||
|
// Create NATS event bus
|
||||||
|
natsBus, err := aether.NewNATSEventBus(nc)
|
||||||
|
require.NoError(t, err)
|
||||||
|
defer natsBus.Stop()
|
||||||
|
|
||||||
|
// Create event store with broadcaster
|
||||||
|
store, err := NewJetStreamEventStoreWithBroadcaster(
|
||||||
|
nc,
|
||||||
|
streamName,
|
||||||
|
natsBus,
|
||||||
|
"tenant-single",
|
||||||
|
)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
// Subscribe to events
|
||||||
|
eventCh := natsBus.Subscribe("tenant-single")
|
||||||
|
|
||||||
|
// Save event
|
||||||
|
testEvent := &aether.Event{
|
||||||
|
ID: "event-1",
|
||||||
|
EventType: "TestEvent",
|
||||||
|
ActorID: "actor-1",
|
||||||
|
Version: 1,
|
||||||
|
Data: map[string]interface{}{"test": "single"},
|
||||||
|
Timestamp: time.Now(),
|
||||||
|
}
|
||||||
|
|
||||||
|
err = store.SaveEvent(testEvent)
|
||||||
|
require.NoError(t, err)
|
||||||
|
log.Printf("Saved event: %s", testEvent.ID)
|
||||||
|
|
||||||
|
// Receive event via NATS (NATSBroadcast may have wrapped in EventStored)
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 1*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
select {
|
||||||
|
case received := <-eventCh:
|
||||||
|
// The received event should have the same actor and version
|
||||||
|
assert.Equal(t, "actor-1", received.ActorID)
|
||||||
|
assert.Equal(t, int64(1), received.Version)
|
||||||
|
// Event type might be original or EventStored wrapper
|
||||||
|
if received.EventType == aether.EventTypeEventStored {
|
||||||
|
assert.Equal(t, "actor-1", received.Data["actorId"].(string))
|
||||||
|
log.Printf("Received EventStored wrapper: %s", received.ID)
|
||||||
|
} else {
|
||||||
|
log.Printf("Received via NATS: %s (type: %s)", received.ID, received.EventType)
|
||||||
|
}
|
||||||
|
case <-ctx.Done():
|
||||||
|
t.Fatal("Did not receive event via NATS")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestCrossNodeBroadcasting_MultiNode tests broadcasting between two nodes.
|
||||||
|
func TestCrossNodeBroadcasting_MultiNode(t *testing.T) {
|
||||||
|
nc, err := nats.Connect(nats.DefaultURL)
|
||||||
|
require.NoError(t, err)
|
||||||
|
defer nc.Close()
|
||||||
|
|
||||||
|
streamName := generateStreamName("broadcast_multi")
|
||||||
|
cleanupStream(nc, streamName)
|
||||||
|
|
||||||
|
// Create Node A
|
||||||
|
nodeANatsBus, err := aether.NewNATSEventBus(nc)
|
||||||
|
require.NoError(t, err)
|
||||||
|
defer nodeANatsBus.Stop()
|
||||||
|
|
||||||
|
nodeAStore, err := NewJetStreamEventStoreWithBroadcaster(
|
||||||
|
nc,
|
||||||
|
streamName,
|
||||||
|
nodeANatsBus,
|
||||||
|
"tenant-multi",
|
||||||
|
)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
// Create Node B
|
||||||
|
nodeBNatsBus, err := aether.NewNATSEventBus(nc)
|
||||||
|
require.NoError(t, err)
|
||||||
|
defer nodeBNatsBus.Stop()
|
||||||
|
|
||||||
|
nodeBStore, err := NewJetStreamEventStoreWithBroadcaster(
|
||||||
|
nc,
|
||||||
|
streamName,
|
||||||
|
nodeBNatsBus,
|
||||||
|
"tenant-multi",
|
||||||
|
)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
// Node A subscribes
|
||||||
|
nodeAEventCh := nodeANatsBus.Subscribe("tenant-multi")
|
||||||
|
|
||||||
|
// Node B subscribes
|
||||||
|
nodeBEventCh := nodeBNatsBus.Subscribe("tenant-multi")
|
||||||
|
|
||||||
|
// Node A saves event
|
||||||
|
eventA := &aether.Event{
|
||||||
|
ID: "event-node-a",
|
||||||
|
EventType: "TestEvent",
|
||||||
|
ActorID: "multi-actor",
|
||||||
|
Version: 1,
|
||||||
|
Data: map[string]interface{}{"node": "a"},
|
||||||
|
Timestamp: time.Now(),
|
||||||
|
}
|
||||||
|
|
||||||
|
err = nodeAStore.SaveEvent(eventA)
|
||||||
|
require.NoError(t, err)
|
||||||
|
log.Printf("Node A saved: %s", eventA.ID)
|
||||||
|
|
||||||
|
// Node B receives event (EventStored wrapper or original)
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 1*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
select {
|
||||||
|
case received := <-nodeBEventCh:
|
||||||
|
// Check actor and version match
|
||||||
|
assert.Equal(t, "multi-actor", received.ActorID)
|
||||||
|
assert.Equal(t, int64(1), received.Version)
|
||||||
|
if received.EventType == aether.EventTypeEventStored {
|
||||||
|
// EventStored wrapper
|
||||||
|
assert.Equal(t, "multi-actor", received.Data["actorId"].(string))
|
||||||
|
} else {
|
||||||
|
// Original event
|
||||||
|
assert.Equal(t, "a", received.Data["node"])
|
||||||
|
}
|
||||||
|
log.Printf("Node B received: %s (actor: %s, version: %d)", received.ID, received.ActorID, received.Version)
|
||||||
|
case <-ctx.Done():
|
||||||
|
t.Fatal("Node B did not receive event")
|
||||||
|
}
|
||||||
|
|
||||||
|
// Give Node B time to receive and process Node A's event
|
||||||
|
time.Sleep(100 * time.Millisecond)
|
||||||
|
|
||||||
|
// Node B saves event with different actor
|
||||||
|
eventB := &aether.Event{
|
||||||
|
ID: "event-node-b",
|
||||||
|
EventType: "TestEvent",
|
||||||
|
ActorID: "multi-actor-b", // Different actor
|
||||||
|
Version: 1,
|
||||||
|
Data: map[string]interface{}{"node": "b"},
|
||||||
|
Timestamp: time.Now(),
|
||||||
|
}
|
||||||
|
|
||||||
|
err = nodeBStore.SaveEvent(eventB)
|
||||||
|
require.NoError(t, err)
|
||||||
|
log.Printf("Node B saved: %s", eventB.ID)
|
||||||
|
|
||||||
|
// Node A receives Node B's event (could be EventStored or original)
|
||||||
|
ctx2, cancel2 := context.WithTimeout(context.Background(), 2*time.Second)
|
||||||
|
defer cancel2()
|
||||||
|
|
||||||
|
// Wait for either EventStored or the actual event
|
||||||
|
received := false
|
||||||
|
for !received {
|
||||||
|
select {
|
||||||
|
case event := <-nodeAEventCh:
|
||||||
|
// Check for Node B's event
|
||||||
|
if event.EventType == aether.EventTypeEventStored {
|
||||||
|
// EventStored wrapper - check actorId
|
||||||
|
actorID := event.Data["actorId"].(string)
|
||||||
|
if actorID == "multi-actor-b" {
|
||||||
|
received = true
|
||||||
|
log.Printf("Node A received EventStored for Node B's event: %s", event.ID)
|
||||||
|
}
|
||||||
|
} else if event.ActorID == "multi-actor-b" {
|
||||||
|
received = true
|
||||||
|
log.Printf("Node A received Node B's event: %s", event.ID)
|
||||||
|
}
|
||||||
|
// Keep listening until we get Node B's event
|
||||||
|
case <-ctx2.Done():
|
||||||
|
if !received {
|
||||||
|
t.Fatal("Node A did not receive Node B's event")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestUpdateVersionCache tests the version cache update logic.
|
||||||
|
func TestUpdateVersionCache(t *testing.T) {
|
||||||
|
nc, err := nats.Connect(nats.DefaultURL)
|
||||||
|
require.NoError(t, err)
|
||||||
|
defer nc.Close()
|
||||||
|
|
||||||
|
streamName := generateStreamName("cache_test")
|
||||||
|
cleanupStream(nc, streamName)
|
||||||
|
|
||||||
|
store, err := NewJetStreamEventStore(nc, streamName)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
// Save event version 1
|
||||||
|
event1 := &aether.Event{
|
||||||
|
ID: "event-1",
|
||||||
|
EventType: "TestEvent",
|
||||||
|
ActorID: "actor-1",
|
||||||
|
Version: 1,
|
||||||
|
Data: map[string]interface{}{},
|
||||||
|
Timestamp: time.Now(),
|
||||||
|
}
|
||||||
|
err = store.SaveEvent(event1)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
// EventStored will trigger UpdateVersionCache
|
||||||
|
time.Sleep(100 * time.Millisecond)
|
||||||
|
|
||||||
|
// Save event version 2 - should succeed (version > 1)
|
||||||
|
event2 := &aether.Event{
|
||||||
|
ID: "event-2",
|
||||||
|
EventType: "TestEvent",
|
||||||
|
ActorID: "actor-1",
|
||||||
|
Version: 2,
|
||||||
|
Data: map[string]interface{}{},
|
||||||
|
Timestamp: time.Now(),
|
||||||
|
}
|
||||||
|
err = store.SaveEvent(event2)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
// Save event version 3 - should succeed (version > 2)
|
||||||
|
event3 := &aether.Event{
|
||||||
|
ID: "event-3",
|
||||||
|
EventType: "TestEvent",
|
||||||
|
ActorID: "actor-1",
|
||||||
|
Version: 3,
|
||||||
|
Data: map[string]interface{}{},
|
||||||
|
Timestamp: time.Now(),
|
||||||
|
}
|
||||||
|
err = store.SaveEvent(event3)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
// Verify all events can be retrieved
|
||||||
|
events, err := store.GetEvents("actor-1", 0)
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Len(t, events, 3)
|
||||||
|
|
||||||
|
// Verify latest version
|
||||||
|
latest, err := store.GetLatestVersion("actor-1")
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Equal(t, int64(3), latest)
|
||||||
|
|
||||||
|
// Manually update cache with version 5 (simulating external update)
|
||||||
|
store.UpdateVersionCache("actor-1", 5)
|
||||||
|
|
||||||
|
// Verify version 4 would conflict (4 < cached 5)
|
||||||
|
event4 := &aether.Event{
|
||||||
|
ID: "event-4",
|
||||||
|
EventType: "TestEvent",
|
||||||
|
ActorID: "actor-1",
|
||||||
|
Version: 4,
|
||||||
|
Data: map[string]interface{}{},
|
||||||
|
Timestamp: time.Now(),
|
||||||
|
}
|
||||||
|
err = store.SaveEvent(event4)
|
||||||
|
assert.Error(t, err, "version 4 should conflict with cached version 5")
|
||||||
|
|
||||||
|
// Version 6 should succeed (6 > 5)
|
||||||
|
event6 := &aether.Event{
|
||||||
|
ID: "event-6",
|
||||||
|
EventType: "TestEvent",
|
||||||
|
ActorID: "actor-1",
|
||||||
|
Version: 6,
|
||||||
|
Data: map[string]interface{}{},
|
||||||
|
Timestamp: time.Now(),
|
||||||
|
}
|
||||||
|
err = store.SaveEvent(event6)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
// Verify version 6 was saved
|
||||||
|
latest, err = store.GetLatestVersion("actor-1")
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Equal(t, int64(6), latest)
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestSubscribeToEventStored tests the convenience helper.
|
||||||
|
func TestSubscribeToEventStored(t *testing.T) {
|
||||||
|
nc, err := nats.Connect(nats.DefaultURL)
|
||||||
|
require.NoError(t, err)
|
||||||
|
defer nc.Close()
|
||||||
|
|
||||||
|
natsBus, err := aether.NewNATSEventBus(nc)
|
||||||
|
require.NoError(t, err)
|
||||||
|
defer natsBus.Stop()
|
||||||
|
|
||||||
|
// Use helper
|
||||||
|
eventStoredCh := natsBus.SubscribeToEventStored("test-store")
|
||||||
|
|
||||||
|
// Verify channel is created
|
||||||
|
assert.NotNil(t, eventStoredCh)
|
||||||
|
|
||||||
|
// Publish EventStored manually (version will be float64 from JSON)
|
||||||
|
eventStored := &aether.Event{
|
||||||
|
ID: "stored-1",
|
||||||
|
EventType: aether.EventTypeEventStored,
|
||||||
|
ActorID: "actor-1",
|
||||||
|
Version: 1,
|
||||||
|
Data: map[string]interface{}{
|
||||||
|
"actorId": "actor-1",
|
||||||
|
"version": 1.0, // Use float64 to match JSON encoding
|
||||||
|
},
|
||||||
|
Timestamp: time.Now(),
|
||||||
|
}
|
||||||
|
|
||||||
|
natsBus.Publish("test-store", eventStored)
|
||||||
|
|
||||||
|
// Should receive the EventStored
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 500*time.Millisecond)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
select {
|
||||||
|
case received := <-eventStoredCh:
|
||||||
|
assert.Equal(t, aether.EventTypeEventStored, received.EventType)
|
||||||
|
assert.Equal(t, "actor-1", received.ActorID)
|
||||||
|
log.Printf("Received EventStored: %s", received.ID)
|
||||||
|
case <-ctx.Done():
|
||||||
|
t.Fatal("Did not receive EventStored")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestCrossNodeBroadcasting_NamespaceIsolation tests namespace isolation.
|
||||||
|
func TestCrossNodeBroadcasting_NamespaceIsolation(t *testing.T) {
|
||||||
|
nc, err := nats.Connect(nats.DefaultURL)
|
||||||
|
require.NoError(t, err)
|
||||||
|
defer nc.Close()
|
||||||
|
|
||||||
|
streamName := generateStreamName("namespace_isolation")
|
||||||
|
cleanupStream(nc, streamName)
|
||||||
|
|
||||||
|
// Create stores with different namespaces
|
||||||
|
storeA, err := NewJetStreamEventStoreWithBroadcaster(
|
||||||
|
nc,
|
||||||
|
streamName,
|
||||||
|
nil,
|
||||||
|
"tenant-a",
|
||||||
|
)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
storeB, err := NewJetStreamEventStoreWithBroadcaster(
|
||||||
|
nc,
|
||||||
|
streamName,
|
||||||
|
nil,
|
||||||
|
"tenant-b",
|
||||||
|
)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
// Save to each namespace
|
||||||
|
eventA := &aether.Event{
|
||||||
|
ID: "event-a",
|
||||||
|
EventType: "TestEvent",
|
||||||
|
ActorID: "actor-a",
|
||||||
|
Version: 1,
|
||||||
|
Data: map[string]interface{}{"tenant": "a"},
|
||||||
|
Timestamp: time.Now(),
|
||||||
|
}
|
||||||
|
|
||||||
|
err = storeA.SaveEvent(eventA)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
eventB := &aether.Event{
|
||||||
|
ID: "event-b",
|
||||||
|
EventType: "TestEvent",
|
||||||
|
ActorID: "actor-b",
|
||||||
|
Version: 1,
|
||||||
|
Data: map[string]interface{}{"tenant": "b"},
|
||||||
|
Timestamp: time.Now(),
|
||||||
|
}
|
||||||
|
|
||||||
|
err = storeB.SaveEvent(eventB)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
// Verify each store can see its own events
|
||||||
|
eventsA, err := storeA.GetEvents("actor-a", 0)
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Len(t, eventsA, 1)
|
||||||
|
assert.Equal(t, "a", eventsA[0].Data["tenant"])
|
||||||
|
|
||||||
|
eventsB, err := storeB.GetEvents("actor-b", 0)
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Len(t, eventsB, 1)
|
||||||
|
assert.Equal(t, "b", eventsB[0].Data["tenant"])
|
||||||
|
}
|
||||||
@@ -558,5 +558,18 @@ func sanitizeSubject(s string) string {
|
|||||||
return s
|
return s
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// UpdateVersionCache updates the version cache for a specific actor.
|
||||||
|
// This is used when receiving events from other nodes via NATS to keep
|
||||||
|
// the version cache consistent across cluster nodes.
|
||||||
|
func (jes *JetStreamEventStore) UpdateVersionCache(actorID string, version int64) {
|
||||||
|
jes.mu.Lock()
|
||||||
|
defer jes.mu.Unlock()
|
||||||
|
|
||||||
|
// Only update if the new version is greater than cached version
|
||||||
|
if currentVersion, ok := jes.versions[actorID]; !ok || version > currentVersion {
|
||||||
|
jes.versions[actorID] = version
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// Compile-time check that JetStreamEventStore implements EventStoreWithErrors
|
// Compile-time check that JetStreamEventStore implements EventStoreWithErrors
|
||||||
var _ aether.EventStoreWithErrors = (*JetStreamEventStore)(nil)
|
var _ aether.EventStoreWithErrors = (*JetStreamEventStore)(nil)
|
||||||
|
|||||||
Reference in New Issue
Block a user