mirror of
https://github.com/garethgeorge/backrest.git
synced 2026-10-05 14:11:25 +00:00
progress implementing remote instance views
This commit is contained in:
1 parent
74eb15df04
commit
ceb29a80f9
7 files changed
+430
-202
No files matched your search
@@ -25,8 +25,8 @@ type PeerState struct {
|
||||
ConnectionStateMessage string
|
||||
|
||||
// Plans and repos available on this peer
|
||||
KnownRepos map[string]struct{}
|
||||
KnownPlans map[string]struct{}
|
||||
KnownRepos map[string]*v1.SyncRepoMetadata
|
||||
KnownPlans map[string]*v1.SyncPlanMetadata
|
||||
|
||||
// Partial configuration available for this peer
|
||||
Config *v1.RemoteConfig
|
||||
@@ -39,8 +39,8 @@ func newPeerState(instanceID, keyID string) *PeerState {
|
||||
LastHeartbeat: time.Now(),
|
||||
ConnectionState: v1.SyncConnectionState_CONNECTION_STATE_DISCONNECTED,
|
||||
ConnectionStateMessage: "disconnected",
|
||||
KnownRepos: make(map[string]struct{}),
|
||||
KnownPlans: make(map[string]struct{}),
|
||||
KnownRepos: make(map[string]*v1.SyncRepoMetadata),
|
||||
KnownPlans: make(map[string]*v1.SyncPlanMetadata),
|
||||
Config: nil, // Will be set when the config is received
|
||||
}
|
||||
}
|
||||
@@ -68,8 +68,8 @@ func peerStateToProto(state *PeerState) *v1.PeerState {
|
||||
LastHeartbeatMillis: state.LastHeartbeat.UnixMilli(),
|
||||
State: state.ConnectionState,
|
||||
StatusMessage: state.ConnectionStateMessage,
|
||||
KnownRepos: slices.Collect(maps.Keys(state.KnownRepos)),
|
||||
KnownPlans: slices.Collect(maps.Keys(state.KnownPlans)),
|
||||
KnownRepos: slices.Collect(maps.Values(state.KnownRepos)),
|
||||
KnownPlans: slices.Collect(maps.Values(state.KnownPlans)),
|
||||
RemoteConfig: state.Config,
|
||||
}
|
||||
}
|
||||
@@ -78,13 +78,13 @@ func peerStateFromProto(state *v1.PeerState) *PeerState {
|
||||
if state.PeerInstanceId == "" || state.PeerKeyid == "" {
|
||||
return nil
|
||||
}
|
||||
knownRepos := make(map[string]struct{}, len(state.KnownRepos))
|
||||
knownRepos := make(map[string]*v1.SyncRepoMetadata, len(state.KnownRepos))
|
||||
for _, repo := range state.KnownRepos {
|
||||
knownRepos[repo] = struct{}{}
|
||||
knownRepos[repo.Id] = repo
|
||||
}
|
||||
knownPlans := make(map[string]struct{}, len(state.KnownPlans))
|
||||
knownPlans := make(map[string]*v1.SyncPlanMetadata, len(state.KnownPlans))
|
||||
for _, plan := range state.KnownPlans {
|
||||
knownPlans[plan] = struct{}{}
|
||||
knownPlans[plan.Id] = plan
|
||||
}
|
||||
|
||||
return &PeerState{
|
||||
@@ -269,20 +269,18 @@ func (m *SqlitePeerStateManager) UpdatePeerState(keyID string, instanceID string
|
||||
defer m.mu.Unlock()
|
||||
|
||||
var state *PeerState
|
||||
stateBytes, err := m.kvstore.Get(keyID)
|
||||
if err != nil {
|
||||
zap.S().Warnf("error getting peer state for key %s: %v", keyID, err)
|
||||
state = newPeerState(instanceID, keyID)
|
||||
} else if stateBytes == nil {
|
||||
state = newPeerState(instanceID, keyID)
|
||||
} else {
|
||||
if stateBytes, err := m.kvstore.Get(keyID); err == nil {
|
||||
var stateProto v1.PeerState
|
||||
if err := proto.Unmarshal(stateBytes, &stateProto); err != nil {
|
||||
zap.S().Warnf("error unmarshalling peer state for key %s: %v", keyID, err)
|
||||
state = newPeerState(instanceID, keyID)
|
||||
} else {
|
||||
state = peerStateFromProto(&stateProto)
|
||||
}
|
||||
} else {
|
||||
zap.S().Warnf("error getting peer state for key %s: %v", keyID, err)
|
||||
}
|
||||
if state == nil {
|
||||
state = newPeerState(instanceID, keyID)
|
||||
}
|
||||
|
||||
updateFn(state)
|
||||
|
||||
@@ -231,13 +231,22 @@ func (c *syncSessionHandlerClient) OnConnectionEstablished(ctx context.Context,
|
||||
for _, repo := range localConfig.Repos {
|
||||
if c.permissions.CheckPermissionForRepo(repo.Guid, v1.Multihost_Permission_PERMISSION_READ_CONFIG) {
|
||||
remoteConfig.Repos = append(remoteConfig.Repos, repo)
|
||||
resourceList.RepoIds = append(resourceList.RepoIds, repo.Id)
|
||||
}
|
||||
if c.permissions.CheckPermissionForRepo(repo.Guid, v1.Multihost_Permission_PERMISSION_READ_OPERATIONS, v1.Multihost_Permission_PERMISSION_READ_CONFIG) {
|
||||
resourceList.Repos = append(resourceList.Repos, &v1.SyncRepoMetadata{
|
||||
Id: repo.Id,
|
||||
Guid: repo.Guid,
|
||||
})
|
||||
}
|
||||
}
|
||||
for _, plan := range localConfig.Plans {
|
||||
if c.permissions.CheckPermissionForPlan(plan.Id, v1.Multihost_Permission_PERMISSION_READ_CONFIG) {
|
||||
remoteConfig.Plans = append(remoteConfig.Plans, plan)
|
||||
resourceList.PlanIds = append(resourceList.PlanIds, plan.Id)
|
||||
}
|
||||
if c.permissions.CheckPermissionForPlan(plan.Id, v1.Multihost_Permission_PERMISSION_READ_OPERATIONS, v1.Multihost_Permission_PERMISSION_READ_CONFIG) {
|
||||
resourceList.Plans = append(resourceList.Plans, &v1.SyncPlanMetadata{
|
||||
Id: plan.Id,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -353,20 +353,20 @@ func (h *syncSessionHandlerServer) HandleSendConfig(ctx context.Context, stream
|
||||
|
||||
func (h *syncSessionHandlerServer) HandleListResources(ctx context.Context, stream *bidiSyncCommandStream, item *v1.SyncStreamItem_SyncActionListResources) error {
|
||||
zap.L().Debug("syncserver received resource list from client", zap.String("client_instance_id", h.peer.InstanceId),
|
||||
zap.Any("repos", item.GetRepoIds()),
|
||||
zap.Any("plans", item.GetPlanIds()))
|
||||
zap.Any("repos", item.GetRepos()),
|
||||
zap.Any("plans", item.GetPlans()))
|
||||
h.mgr.peerStateManager.UpdatePeerState(h.peer.Keyid, h.peer.InstanceId, func(peerState *PeerState) {
|
||||
if peerState == nil {
|
||||
return // this should not happen
|
||||
}
|
||||
|
||||
repos := item.GetRepoIds()
|
||||
plans := item.GetPlanIds()
|
||||
for _, repoID := range repos {
|
||||
peerState.KnownRepos[repoID] = struct{}{}
|
||||
repos := item.GetRepos()
|
||||
plans := item.GetPlans()
|
||||
for _, repo := range repos {
|
||||
peerState.KnownRepos[repo.Id] = repo
|
||||
}
|
||||
for _, planID := range plans {
|
||||
peerState.KnownPlans[planID] = struct{}{}
|
||||
for _, plan := range plans {
|
||||
peerState.KnownPlans[plan.Id] = plan
|
||||
}
|
||||
})
|
||||
return nil
|
||||
|
||||
Reference in new issue
Block a user