From 134e5324d7c33d6c38160204a499ccc43f8cb753 Mon Sep 17 00:00:00 2001 From: Dmitrii Stepanov Date: Wed, 19 Feb 2025 10:30:48 +0300 Subject: [PATCH] [#0] tree: Split GetOpLog stream Signed-off-by: Dmitrii Stepanov --- pkg/services/tree/sync.go | 71 +++++++++++++++++++++++++-------------- 1 file changed, 45 insertions(+), 26 deletions(-) diff --git a/pkg/services/tree/sync.go b/pkg/services/tree/sync.go index 0f85f50b1..45b62c51a 100644 --- a/pkg/services/tree/sync.go +++ b/pkg/services/tree/sync.go @@ -219,38 +219,57 @@ func (s *Service) startStream(ctx context.Context, cid cid.ID, treeID string, rawCID := make([]byte, sha256.Size) cid.Encode(rawCID) + from := height + const batchSize = 10_000 - req := &GetOpLogRequest{ - Body: &GetOpLogRequest_Body{ - ContainerId: rawCID, - TreeId: treeID, - Height: height, - }, - } - if err := SignMessage(req, s.key); err != nil { - return err - } - - c, err := treeClient.GetOpLog(ctx, req) - if err != nil { - return fmt.Errorf("can't initialize client: %w", err) - } - res, err := c.Recv() - for ; err == nil; res, err = c.Recv() { - lm := res.GetBody().GetOperation() - m := &pilorama.Move{ - Parent: lm.GetParentId(), - Child: lm.GetChildId(), + for { + count := 0 + req := &GetOpLogRequest{ + Body: &GetOpLogRequest_Body{ + ContainerId: rawCID, + TreeId: treeID, + Height: from, + }, } - if err := m.Meta.FromBytes(lm.GetMeta()); err != nil { + if err := SignMessage(req, s.key); err != nil { return err } - opsCh <- m - } - if !errors.Is(err, io.EOF) { + + streamCtx, cancel := context.WithCancel(ctx) + defer cancel() + + c, err := treeClient.GetOpLog(streamCtx, req) + if err != nil { + return fmt.Errorf("can't initialize client: %w", err) + } + res, err := c.Recv() + for ; err == nil; res, err = c.Recv() { + lm := res.GetBody().GetOperation() + m := &pilorama.Move{ + Parent: lm.GetParentId(), + Child: lm.GetChildId(), + } + if err := m.Meta.FromBytes(lm.GetMeta()); err != nil { + return err + } + from = m.Time + opsCh <- m + count++ + if count == batchSize { + break // c.Recv() + } + } + if err == nil { + // count == batchSize + // close current stream and start new one + cancel() + continue + } + if errors.Is(err, io.EOF) { + return nil + } return err } - return nil } // synchronizeTree synchronizes operations getting them from different nodes.