From bdbd2b46c0c369c5e4c6136e891dafa9d3eac187 Mon Sep 17 00:00:00 2001 From: Gareth George Date: Tue, 21 Apr 2026 23:54:37 -0700 Subject: [PATCH] logging improvements --- internal/api/syncapi/syncclient.go | 93 +++++++++++++++++++++--------- internal/api/syncapi/syncserver.go | 41 +++++++------ 2 files changed, 89 insertions(+), 45 deletions(-) diff --git a/internal/api/syncapi/syncclient.go b/internal/api/syncapi/syncclient.go index 8b864b58..d026b9b2 100644 --- a/internal/api/syncapi/syncclient.go +++ b/internal/api/syncapi/syncclient.go @@ -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 } diff --git a/internal/api/syncapi/syncserver.go b/internal/api/syncapi/syncserver.go index a8747474..cb036824 100644 --- a/internal/api/syncapi/syncserver.go +++ b/internal/api/syncapi/syncserver.go @@ -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)