retryOperation attempts up to `count` attempts at running given function, exiting as soon as it returns with non-error. gno: only for logging
(subject string, txMsg []byte, gno int64)
| 778 | // exiting as soon as it returns with non-error. |
| 779 | // gno: only for logging |
| 780 | func (e *ExtractorOracle) publish(subject string, txMsg []byte, gno int64) (err error) { |
| 781 | msgLen := len(txMsg) |
| 782 | |
| 783 | data := txMsg |
| 784 | lenData := len(data) |
| 785 | |
| 786 | // lenData < NatsMaxMsg: 1 msg |
| 787 | // lenData = k * NatsMaxMsg + b, where k >= 1 && b >= 0: (k+1) msg |
| 788 | // b could be 0. we send a zero-len msg as a sign of termination. |
| 789 | nSeg := lenData/g.NatsMaxMsg + 1 |
| 790 | e.logger.Debug("publish. msg", "subject", subject, "gno", gno, "nSeg", nSeg, "spanLen", lenData, "msgLen", msgLen) |
| 791 | bak := make([]byte, 4) |
| 792 | if nSeg > 1 { |
| 793 | // ensure there are 4 bytes to save iSeg |
| 794 | data = append(data, 0, 0, 0, 0) |
| 795 | } |
| 796 | for iSeg := 0; iSeg < nSeg; iSeg++ { |
| 797 | var part []byte |
| 798 | if nSeg == 1 { // not big msg |
| 799 | part = data |
| 800 | } else { |
| 801 | begin := iSeg * g.NatsMaxMsg |
| 802 | end := mathutil.Min(lenData, (iSeg+1)*g.NatsMaxMsg) |
| 803 | // use extra 4 bytes to save iSeg |
| 804 | if iSeg > 0 { |
| 805 | copy(data[begin:begin+4], bak) |
| 806 | } |
| 807 | copy(bak, data[end:end+4]) |
| 808 | part = data[begin : end+4] |
| 809 | binary.LittleEndian.PutUint32(data[end:], uint32(iSeg)) |
| 810 | } |
| 811 | |
| 812 | e.logger.Debug("publish", "subject", subject, "gno", gno, "partLen", len(part), "iSeg", iSeg) |
| 813 | _, err := e.natsConn.Request(subject, part, 24*time.Hour) |
| 814 | if err != nil { |
| 815 | e.logger.Error("unexpected error on publish", "err", err) |
| 816 | return err |
| 817 | } |
| 818 | } |
| 819 | return nil |
| 820 | } |
| 821 | |
| 822 | func (e *ExtractorOracle) onError(state int, err error) { |
| 823 | e.logger.Error("onError", "err", err, "hasShutdown", e.shutdown) |
no outgoing calls
no test coverage detected