mirror of
https://github.com/netbirdio/netbird.git
synced 2026-04-18 08:16:39 +00:00
fix bug with stream
This commit is contained in:
@@ -898,7 +898,6 @@ func (e *Engine) receiveJobEvents() {
|
|||||||
bundleResult, err := e.handleBundle(params.Bundle)
|
bundleResult, err := e.handleBundle(params.Bundle)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
resp.Reason = []byte(err.Error())
|
resp.Reason = []byte(err.Error())
|
||||||
resp.Status = mgmProto.JobStatus_failed
|
|
||||||
return &resp
|
return &resp
|
||||||
}
|
}
|
||||||
resp.Status = mgmProto.JobStatus_succeeded
|
resp.Status = mgmProto.JobStatus_succeeded
|
||||||
|
|||||||
@@ -160,6 +160,8 @@ func (j *Job) ApplyResponse(resp *proto.JobResponse) error {
|
|||||||
if resp == nil {
|
if resp == nil {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
j.ID = string(resp.ID)
|
||||||
now := time.Now().UTC()
|
now := time.Now().UTC()
|
||||||
j.CompletedAt = &now
|
j.CompletedAt = &now
|
||||||
switch resp.Status {
|
switch resp.Status {
|
||||||
|
|||||||
@@ -180,13 +180,19 @@ func (c *GrpcClient) handleJobStream(
|
|||||||
// Main loop: receive, process, respond
|
// Main loop: receive, process, respond
|
||||||
for {
|
for {
|
||||||
jobReq, err := c.receiveJobRequest(ctx, stream, serverPubKey)
|
jobReq, err := c.receiveJobRequest(ctx, stream, serverPubKey)
|
||||||
if err != nil {
|
if err != nil && err != io.EOF {
|
||||||
if errors.Is(err, io.EOF) || errors.Is(err, context.Canceled) {
|
c.notifyDisconnected(err)
|
||||||
log.WithContext(ctx).Info("job stream closed by server")
|
s, _ := gstatus.FromError(err)
|
||||||
|
switch s.Code() {
|
||||||
|
case codes.PermissionDenied:
|
||||||
|
return backoff.Permanent(err) // unrecoverable error, propagate to the upper layer
|
||||||
|
case codes.Canceled:
|
||||||
|
log.Debugf("management connection context has been canceled, this usually indicates shutdown")
|
||||||
return nil
|
return nil
|
||||||
|
default:
|
||||||
|
log.Warnf("disconnected from the Management service but will retry silently. Reason: %v", err)
|
||||||
|
return err
|
||||||
}
|
}
|
||||||
log.WithContext(ctx).Errorf("error receiving job request: %v", err)
|
|
||||||
return err
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if jobReq == nil || len(jobReq.ID) == 0 {
|
if jobReq == nil || len(jobReq.ID) == 0 {
|
||||||
@@ -298,7 +304,7 @@ func (c *GrpcClient) handleSyncStream(ctx context.Context, serverPubKey wgtypes.
|
|||||||
|
|
||||||
// blocking until error
|
// blocking until error
|
||||||
err = c.receiveUpdatesEvents(stream, serverPubKey, msgHandler)
|
err = c.receiveUpdatesEvents(stream, serverPubKey, msgHandler)
|
||||||
if err != nil && err != io.EOF{
|
if err != nil && err != io.EOF {
|
||||||
c.notifyDisconnected(err)
|
c.notifyDisconnected(err)
|
||||||
s, _ := gstatus.FromError(err)
|
s, _ := gstatus.FromError(err)
|
||||||
switch s.Code() {
|
switch s.Code() {
|
||||||
|
|||||||
Reference in New Issue
Block a user