mirror of
https://github.com/garethgeorge/backrest.git
synced 2026-09-29 11:25:41 +00:00
logging improvements
This commit is contained in:
@@ -234,7 +234,7 @@ func (c *syncSessionHandlerClient) canForwardMeta(meta oplog.OpMetadata) bool {
|
||||
return false
|
||||
}
|
||||
|
||||
func (c *syncSessionHandlerClient) sendManifest(stream *bidiSyncCommandStream) error {
|
||||
func (c *syncSessionHandlerClient) sendManifest(stream *bidiSyncCommandStream) (int, error) {
|
||||
var opIDs, modnos []int64
|
||||
if err := c.oplog.QueryMetadata(oplog.Query{}, func(meta oplog.OpMetadata) error {
|
||||
if c.canForwardMeta(meta) {
|
||||
@@ -243,9 +243,8 @@ func (c *syncSessionHandlerClient) sendManifest(stream *bidiSyncCommandStream) e
|
||||
}
|
||||
return nil
|
||||
}); err != nil {
|
||||
return fmt.Errorf("querying operation metadata for manifest: %w", err)
|
||||
return 0, fmt.Errorf("querying operation metadata for manifest: %w", err)
|
||||
}
|
||||
c.l.Sugar().Debugf("sending operation manifest with %d operations", len(opIDs))
|
||||
stream.Send(&v1sync.SyncStreamItem{
|
||||
Action: &v1sync.SyncStreamItem_OperationManifest{
|
||||
OperationManifest: &v1sync.SyncStreamItem_SyncActionOperationManifest{
|
||||
@@ -254,7 +253,7 @@ func (c *syncSessionHandlerClient) sendManifest(stream *bidiSyncCommandStream) e
|
||||
},
|
||||
},
|
||||
})
|
||||
return nil
|
||||
return len(opIDs), nil
|
||||
}
|
||||
|
||||
func (c *syncSessionHandlerClient) OnConnectionEstablished(ctx context.Context, stream *bidiSyncCommandStream, peer *v1.Multihost_Peer) error {
|
||||
@@ -308,12 +307,14 @@ func (c *syncSessionHandlerClient) OnConnectionEstablished(ctx context.Context,
|
||||
go sendHeartbeats(ctx, stream, env.MultihostHeartbeatInterval())
|
||||
|
||||
// Forward a view of our config (if the peer is allowed to see it).
|
||||
if err := c.sendConfig(ctx, stream); err != nil {
|
||||
repoCount, planCount, err := c.sendConfig(ctx, stream)
|
||||
if err != nil {
|
||||
return fmt.Errorf("send config to peer %q: %w", peer.InstanceId, err)
|
||||
}
|
||||
|
||||
// Forward a list of the resources we're making available to the peer
|
||||
if err := c.sendResourceList(ctx, stream); err != nil {
|
||||
resRepoCount, resPlanCount, err := c.sendResourceList(ctx, stream)
|
||||
if err != nil {
|
||||
return fmt.Errorf("send resource list to peer %q: %w", peer.InstanceId, err)
|
||||
}
|
||||
|
||||
@@ -374,15 +375,20 @@ func (c *syncSessionHandlerClient) OnConnectionEstablished(ctx context.Context,
|
||||
}()
|
||||
|
||||
// Send initial operation manifest to the server for reconciliation.
|
||||
if err := c.sendManifest(stream); err != nil {
|
||||
opCount, err := c.sendManifest(stream)
|
||||
if err != nil {
|
||||
return fmt.Errorf("send manifest to peer %q: %w", peer.InstanceId, err)
|
||||
}
|
||||
|
||||
c.l.Sugar().Infof("sent initial state to server: %d operations, %d repos, %d plans (config); %d repos, %d plans (resources)",
|
||||
opCount, repoCount, planCount, resRepoCount, resPlanCount)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *syncSessionHandlerClient) HandleRequestResources(ctx context.Context, stream *bidiSyncCommandStream, item *v1sync.SyncStreamItem_SyncActionRequestResources) error {
|
||||
return c.sendResourceList(ctx, stream)
|
||||
_, _, err := c.sendResourceList(ctx, stream)
|
||||
return err
|
||||
}
|
||||
|
||||
func (c *syncSessionHandlerClient) HandleHeartbeat(ctx context.Context, stream *bidiSyncCommandStream, item *v1sync.SyncStreamItem_SyncActionHeartbeat) error {
|
||||
@@ -397,7 +403,12 @@ func (c *syncSessionHandlerClient) HandleHeartbeat(ctx context.Context, stream *
|
||||
|
||||
func (c *syncSessionHandlerClient) HandleOperationManifest(ctx context.Context, stream *bidiSyncCommandStream, item *v1sync.SyncStreamItem_SyncActionOperationManifest) error {
|
||||
// Server re-requested a manifest (e.g. after reconnect). Respond with a fresh one.
|
||||
return c.sendManifest(stream)
|
||||
opCount, err := c.sendManifest(stream)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
c.l.Sugar().Debugf("re-sent operation manifest with %d operations", opCount)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *syncSessionHandlerClient) HandleRequestOperationData(ctx context.Context, stream *bidiSyncCommandStream, item *v1sync.SyncStreamItem_SyncActionRequestOperationData) error {
|
||||
@@ -438,9 +449,7 @@ func (c *syncSessionHandlerClient) HandleReceiveOperations(ctx context.Context,
|
||||
}
|
||||
|
||||
func (c *syncSessionHandlerClient) HandleReceiveResources(ctx context.Context, stream *bidiSyncCommandStream, item *v1sync.SyncStreamItem_SyncActionReceiveResources) error {
|
||||
c.l.Debug("received resource list from server",
|
||||
zap.Any("repos", item.GetRepos()),
|
||||
zap.Any("plans", item.GetPlans()))
|
||||
c.l.Sugar().Debugf("received resource list from server: %d repos, %d plans", len(item.GetRepos()), len(item.GetPlans()))
|
||||
peerState := c.mgr.peerStateManager.GetPeerState(c.peer.Keyid).Clone()
|
||||
if peerState == nil {
|
||||
return NewSyncErrorInternal(fmt.Errorf("peer state for %q not found", c.peer.Keyid))
|
||||
@@ -459,7 +468,8 @@ func (c *syncSessionHandlerClient) HandleReceiveResources(ctx context.Context, s
|
||||
|
||||
// Note unused: there isn't a situation where the host would send its config for information, the host will only call 'SetConfig' to update the config.
|
||||
func (c *syncSessionHandlerClient) HandleReceiveConfig(ctx context.Context, stream *bidiSyncCommandStream, item *v1sync.SyncStreamItem_SyncActionReceiveConfig) error {
|
||||
c.l.Sugar().Debugf("received remote config update")
|
||||
c.l.Sugar().Debugf("received remote config: %d repos, %d plans, modno=%d",
|
||||
len(item.GetConfig().GetRepos()), len(item.GetConfig().GetPlans()), item.GetConfig().GetModno())
|
||||
peerState := c.mgr.peerStateManager.GetPeerState(c.peer.Keyid).Clone()
|
||||
if peerState == nil {
|
||||
return NewSyncErrorInternal(fmt.Errorf("peer state for %q not found", c.peer.Keyid))
|
||||
@@ -474,13 +484,11 @@ func (c *syncSessionHandlerClient) HandleReceiveConfig(ctx context.Context, stre
|
||||
}
|
||||
|
||||
func (c *syncSessionHandlerClient) HandleSetConfig(ctx context.Context, stream *bidiSyncCommandStream, item *v1sync.SyncStreamItem_SyncActionSetConfig) error {
|
||||
c.l.Sugar().Debugf("received SetConfig request from peer %q", c.peer.GetInstanceId())
|
||||
|
||||
return c.mgr.configMgr.Transform(func(cfg *v1.Config) (*v1.Config, error) {
|
||||
snapshot := proto.Clone(cfg).(*v1.Config) // snapshot for change detection
|
||||
|
||||
var plansNew, plansUpdated, plansUnchanged int
|
||||
for _, plan := range item.GetPlans() {
|
||||
c.l.Sugar().Debugf("received plan update: %s", plan.Id)
|
||||
if !c.permissions.CheckPermissionForPlan(plan.Id, permissions.PermsCanWriteConfiguration...) {
|
||||
return nil, NewSyncErrorAuth(fmt.Errorf("peer %q is not allowed to update plan %q", c.peer.InstanceId, plan.Id))
|
||||
}
|
||||
@@ -489,14 +497,23 @@ func (c *syncSessionHandlerClient) HandleSetConfig(ctx context.Context, stream *
|
||||
return p.Id == plan.Id
|
||||
})
|
||||
if idx >= 0 {
|
||||
if proto.Equal(cfg.Plans[idx], plan) {
|
||||
c.l.Sugar().Debugf("received plan %s (unchanged)", plan.Id)
|
||||
plansUnchanged++
|
||||
} else {
|
||||
c.l.Sugar().Debugf("received plan %s (updated)", plan.Id)
|
||||
plansUpdated++
|
||||
}
|
||||
cfg.Plans[idx] = plan
|
||||
} else {
|
||||
c.l.Sugar().Debugf("received plan %s (new)", plan.Id)
|
||||
plansNew++
|
||||
cfg.Plans = append(cfg.Plans, plan)
|
||||
}
|
||||
}
|
||||
|
||||
var reposNew, reposUpdated, reposUnchanged, reposSkipped int
|
||||
for _, repo := range item.GetRepos() {
|
||||
c.l.Sugar().Debugf("received repo update: %s", repo.Guid)
|
||||
isFromOriginPeer := repo.GetOriginInstanceId() != "" && repo.GetOriginInstanceId() == c.peer.InstanceId
|
||||
if !isFromOriginPeer && !c.permissions.CheckPermissionForRepo(repo.Id, permissions.PermsCanWriteConfiguration...) {
|
||||
return nil, NewSyncErrorAuth(fmt.Errorf("peer %q is not allowed to update repo %q", c.peer.InstanceId, repo.Id))
|
||||
@@ -506,6 +523,13 @@ func (c *syncSessionHandlerClient) HandleSetConfig(ctx context.Context, stream *
|
||||
return r.Guid == repo.Guid
|
||||
})
|
||||
if idx >= 0 {
|
||||
if proto.Equal(cfg.Repos[idx], repo) {
|
||||
c.l.Sugar().Debugf("received repo %s (unchanged)", repo.Id)
|
||||
reposUnchanged++
|
||||
} else {
|
||||
c.l.Sugar().Debugf("received repo %s (updated)", repo.Id)
|
||||
reposUpdated++
|
||||
}
|
||||
cfg.Repos[idx] = repo
|
||||
} else {
|
||||
conflictIdx := slices.IndexFunc(cfg.Repos, func(r *v1.Repo) bool {
|
||||
@@ -513,14 +537,17 @@ func (c *syncSessionHandlerClient) HandleSetConfig(ctx context.Context, stream *
|
||||
})
|
||||
if conflictIdx >= 0 {
|
||||
c.l.Sugar().Warnf("received shared repo %q (guid %s) conflicts with existing local repo %q (guid %s), skipping", repo.Id, repo.Guid, cfg.Repos[conflictIdx].Id, cfg.Repos[conflictIdx].Guid)
|
||||
reposSkipped++
|
||||
continue
|
||||
}
|
||||
c.l.Sugar().Debugf("received repo %s (new)", repo.Id)
|
||||
reposNew++
|
||||
cfg.Repos = append(cfg.Repos, repo)
|
||||
}
|
||||
}
|
||||
|
||||
var plansDeleted int
|
||||
for _, plan := range item.GetPlansToDelete() {
|
||||
c.l.Sugar().Debugf("received plan deletion request: %s", plan)
|
||||
if !c.permissions.CheckPermissionForPlan(plan, permissions.PermsCanWriteConfiguration...) {
|
||||
return nil, NewSyncErrorAuth(fmt.Errorf("peer %q is not allowed to delete plan %q", c.peer.InstanceId, plan))
|
||||
}
|
||||
@@ -529,14 +556,16 @@ func (c *syncSessionHandlerClient) HandleSetConfig(ctx context.Context, stream *
|
||||
return p.Id == plan
|
||||
})
|
||||
if idx >= 0 {
|
||||
c.l.Sugar().Debugf("received plan deletion: %s", plan)
|
||||
plansDeleted++
|
||||
cfg.Plans = append(cfg.Plans[:idx], cfg.Plans[idx+1:]...)
|
||||
} else {
|
||||
c.l.Sugar().Warnf("received plan deletion request for non-existent plan %q, ignoring", plan)
|
||||
}
|
||||
}
|
||||
|
||||
var reposDeleted int
|
||||
for _, repoID := range item.GetReposToDelete() {
|
||||
c.l.Sugar().Debugf("received repo deletion request: %s", repoID)
|
||||
if !c.permissions.CheckPermissionForRepo(repoID, permissions.PermsCanWriteConfiguration...) {
|
||||
return nil, NewSyncErrorAuth(fmt.Errorf("peer %q is not allowed to delete repo %q", c.peer.InstanceId, repoID))
|
||||
}
|
||||
@@ -545,6 +574,8 @@ func (c *syncSessionHandlerClient) HandleSetConfig(ctx context.Context, stream *
|
||||
return r.Id == repoID
|
||||
})
|
||||
if idx >= 0 {
|
||||
c.l.Sugar().Debugf("received repo deletion: %s", repoID)
|
||||
reposDeleted++
|
||||
cfg.Repos = append(cfg.Repos[:idx], cfg.Repos[idx+1:]...)
|
||||
} else {
|
||||
c.l.Sugar().Warnf("received repo deletion request for non-existent repo %q, ignoring", repoID)
|
||||
@@ -552,17 +583,25 @@ func (c *syncSessionHandlerClient) HandleSetConfig(ctx context.Context, stream *
|
||||
}
|
||||
|
||||
// Skip the update if nothing actually changed to avoid triggering a reconnect loop.
|
||||
if proto.Equal(cfg, snapshot) {
|
||||
c.l.Sugar().Debugf("SetConfig resulted in no changes, skipping update")
|
||||
return nil, nil
|
||||
hasChanges := !proto.Equal(cfg, snapshot)
|
||||
if hasChanges {
|
||||
cfg.Modno++
|
||||
}
|
||||
|
||||
cfg.Modno++
|
||||
c.l.Sugar().Debugf("SetConfig from peer %q: repos(%d new, %d updated, %d unchanged, %d skipped, %d deleted) plans(%d new, %d updated, %d unchanged, %d deleted) — config %s",
|
||||
c.peer.GetInstanceId(),
|
||||
reposNew, reposUpdated, reposUnchanged, reposSkipped, reposDeleted,
|
||||
plansNew, plansUpdated, plansUnchanged, plansDeleted,
|
||||
map[bool]string{true: "updated", false: "unchanged"}[hasChanges])
|
||||
|
||||
if !hasChanges {
|
||||
return nil, nil
|
||||
}
|
||||
return cfg, nil
|
||||
})
|
||||
}
|
||||
|
||||
func (c *syncSessionHandlerClient) sendConfig(ctx context.Context, stream *bidiSyncCommandStream) error {
|
||||
func (c *syncSessionHandlerClient) sendConfig(ctx context.Context, stream *bidiSyncCommandStream) (int, int, error) {
|
||||
localConfig := c.syncConfigSnapshot.config
|
||||
remoteConfig := &v1sync.RemoteConfig{
|
||||
Version: localConfig.Version,
|
||||
@@ -589,10 +628,10 @@ func (c *syncSessionHandlerClient) sendConfig(ctx context.Context, stream *bidiS
|
||||
},
|
||||
})
|
||||
|
||||
return nil
|
||||
return len(remoteConfig.Repos), len(remoteConfig.Plans), nil
|
||||
}
|
||||
|
||||
func (c *syncSessionHandlerClient) sendResourceList(ctx context.Context, stream *bidiSyncCommandStream) error {
|
||||
func (c *syncSessionHandlerClient) sendResourceList(ctx context.Context, stream *bidiSyncCommandStream) (int, int, error) {
|
||||
repoMetadatas := []*v1sync.RepoMetadata{}
|
||||
planMetadatas := []*v1sync.PlanMetadata{}
|
||||
|
||||
@@ -622,5 +661,5 @@ func (c *syncSessionHandlerClient) sendResourceList(ctx context.Context, stream
|
||||
},
|
||||
})
|
||||
|
||||
return nil
|
||||
return len(repoMetadatas), len(planMetadatas), nil
|
||||
}
|
||||
|
||||
@@ -196,23 +196,30 @@ func (h *syncSessionHandlerServer) OnConnectionEstablished(ctx context.Context,
|
||||
}
|
||||
|
||||
// Permissions unchanged — send updated config and shared repos to client
|
||||
h.l.Sugar().Debugf("config changed, sending updated config to client %q", h.peer.InstanceId)
|
||||
if err := h.sendConfigToClient(stream, newConfig); err != nil {
|
||||
configRepos, configPlans, err := h.sendConfigToClient(stream, newConfig)
|
||||
if err != nil {
|
||||
h.l.Sugar().Warnf("failed to send updated config to client %q: %v", h.peer.InstanceId, err)
|
||||
} else {
|
||||
sharedRepos := h.sendSharedReposToClient(stream, newConfig)
|
||||
h.l.Sugar().Debugf("config changed, sent update to client %q: %d repos, %d plans (config); %d shared repos pushed",
|
||||
h.peer.InstanceId, configRepos, configPlans, sharedRepos)
|
||||
}
|
||||
h.sendSharedReposToClient(stream, newConfig)
|
||||
case <-ctx.Done():
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
if err := h.sendConfigToClient(stream, h.snapshot.config); err != nil {
|
||||
configRepos, configPlans, err := h.sendConfigToClient(stream, h.snapshot.config)
|
||||
if err != nil {
|
||||
return NewSyncErrorInternal(fmt.Errorf("sending initial config to client: %w", err))
|
||||
}
|
||||
|
||||
// Push shared repos to the client
|
||||
h.sendSharedReposToClient(stream, h.snapshot.config)
|
||||
sharedRepoCount := h.sendSharedReposToClient(stream, h.snapshot.config)
|
||||
|
||||
h.l.Sugar().Infof("sent initial state to client %q: %d repos, %d plans (config); %d shared repos pushed",
|
||||
h.peer.InstanceId, configRepos, configPlans, sharedRepoCount)
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -323,14 +330,12 @@ func (h *syncSessionHandlerServer) deleteByOriginalID(originalID int64) error {
|
||||
return h.mgr.oplog.Delete(foundOp.ID)
|
||||
}
|
||||
|
||||
func (h *syncSessionHandlerServer) sendConfigToClient(stream *bidiSyncCommandStream, config *v1.Config) error {
|
||||
func (h *syncSessionHandlerServer) sendConfigToClient(stream *bidiSyncCommandStream, config *v1.Config) (int, int, error) {
|
||||
remoteConfig := &v1sync.RemoteConfig{
|
||||
Version: config.Version,
|
||||
Modno: config.Modno,
|
||||
}
|
||||
resourceListMsg := &v1sync.SyncStreamItem_SyncActionReceiveResources{}
|
||||
var allowedRepoIDs []string
|
||||
var allowedPlanIDs []string
|
||||
for _, repo := range config.Repos {
|
||||
if h.permissions.CheckPermissionForRepo(repo.Id, permissions.PermsCanViewConfiguration...) {
|
||||
remoteConfig.Repos = append(remoteConfig.Repos, repo)
|
||||
@@ -338,7 +343,6 @@ func (h *syncSessionHandlerServer) sendConfigToClient(stream *bidiSyncCommandStr
|
||||
Id: repo.Id,
|
||||
Guid: repo.Guid,
|
||||
})
|
||||
allowedRepoIDs = append(allowedRepoIDs, repo.Id)
|
||||
}
|
||||
}
|
||||
for _, plan := range config.Plans {
|
||||
@@ -347,10 +351,8 @@ func (h *syncSessionHandlerServer) sendConfigToClient(stream *bidiSyncCommandStr
|
||||
resourceListMsg.Plans = append(resourceListMsg.Plans, &v1sync.PlanMetadata{
|
||||
Id: plan.Id,
|
||||
})
|
||||
allowedPlanIDs = append(allowedPlanIDs, plan.Id)
|
||||
}
|
||||
}
|
||||
h.l.Sugar().Debugf("determined client %v is allowlisted to read configs for repos %v and plans %v", h.peer.InstanceId, allowedRepoIDs, allowedPlanIDs)
|
||||
|
||||
// Send the config, this is the first meaningful packet the client will receive.
|
||||
stream.Send(&v1sync.SyncStreamItem{
|
||||
@@ -367,14 +369,15 @@ func (h *syncSessionHandlerServer) sendConfigToClient(stream *bidiSyncCommandStr
|
||||
ReceiveResources: resourceListMsg,
|
||||
},
|
||||
})
|
||||
return nil
|
||||
return len(remoteConfig.Repos), len(remoteConfig.Plans), nil
|
||||
}
|
||||
|
||||
// sendSharedReposToClient sends repos marked as shared to the client via SetConfig.
|
||||
// This pushes repo configurations to the client so they are added to the client's local config.
|
||||
func (h *syncSessionHandlerServer) sendSharedReposToClient(stream *bidiSyncCommandStream, config *v1.Config) {
|
||||
// Returns the number of shared repos sent.
|
||||
func (h *syncSessionHandlerServer) sendSharedReposToClient(stream *bidiSyncCommandStream, config *v1.Config) int {
|
||||
if !h.permissions.HasPermissionType(v1.Multihost_Permission_PERMISSION_RECEIVE_SHARED_REPOS) {
|
||||
return
|
||||
return 0
|
||||
}
|
||||
|
||||
var sharedRepos []*v1.Repo
|
||||
@@ -387,10 +390,9 @@ func (h *syncSessionHandlerServer) sendSharedReposToClient(stream *bidiSyncComma
|
||||
}
|
||||
|
||||
if len(sharedRepos) == 0 {
|
||||
return
|
||||
return 0
|
||||
}
|
||||
|
||||
h.l.Sugar().Debugf("sending %d shared repos to client %q", len(sharedRepos), h.peer.InstanceId)
|
||||
stream.Send(&v1sync.SyncStreamItem{
|
||||
Action: &v1sync.SyncStreamItem_SetConfig{
|
||||
SetConfig: &v1sync.SyncStreamItem_SyncActionSetConfig{
|
||||
@@ -398,6 +400,7 @@ func (h *syncSessionHandlerServer) sendSharedReposToClient(stream *bidiSyncComma
|
||||
},
|
||||
},
|
||||
})
|
||||
return len(sharedRepos)
|
||||
}
|
||||
|
||||
// ValidatePairingSecret checks a pairing secret against a list of pairing tokens.
|
||||
@@ -522,9 +525,11 @@ func (h *syncSessionHandlerServer) HandleOperationManifest(ctx context.Context,
|
||||
}
|
||||
|
||||
// Find ops we need (new or changed modno), preserving manifest order
|
||||
opIDs := item.GetOpIds()
|
||||
modnos := item.GetModnos()
|
||||
var needIDs []int64
|
||||
for i, id := range item.GetOpIds() {
|
||||
modno := item.GetModnos()[i]
|
||||
for i, id := range opIDs {
|
||||
modno := modnos[i]
|
||||
local, exists := localState[id]
|
||||
if !exists || local.modno != modno {
|
||||
needIDs = append(needIDs, id)
|
||||
|
||||
Reference in New Issue
Block a user