File size: 1,987 Bytes
93d826e
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
package api

import (
	"bufio"
	"bytes"
	"encoding/json"
	"fmt"
	"io"
	"log"
	"plandex-cli/types"
	"time"

	shared "plandex-shared"
)

// 3 heartbeat misses = timeout
const HeartbeatTimeout = 16 * time.Second

func connectPlanRespStream(body io.ReadCloser, onStream types.OnStreamPlan) {
	reader := bufio.NewReader(body)
	timer := time.NewTimer(HeartbeatTimeout)
	defer timer.Stop()

	go func() {
		for {
			select {
			case <-timer.C:
				log.Println("Connection to plan stream timed out due to missing heartbeats")
				onStream(types.OnStreamPlanParams{Msg: nil, Err: fmt.Errorf("connection to plan stream timed out due to missing heartbeats")})
				body.Close()
				return
			default:
			}

			s, err := readUntilSeparator(reader, shared.STREAM_MESSAGE_SEPARATOR)
			if err != nil {
				log.Println("Error reading line:", err)
				onStream(types.OnStreamPlanParams{Msg: nil, Err: err})
				body.Close()
				return
			}

			timer.Reset(HeartbeatTimeout)

			// ignore heartbeats
			if s == string(shared.StreamMessageHeartbeat) {
				continue
			}

			var msg shared.StreamMessage
			err = json.Unmarshal([]byte(s), &msg)
			if err != nil {
				log.Println("Error unmarshalling message:", err)
				onStream(types.OnStreamPlanParams{Msg: nil, Err: err})
				body.Close()
				return
			}

			// log.Println("connectPlanRespStream: received message:", msg)

			onStream(types.OnStreamPlanParams{Msg: &msg, Err: nil})

			if msg.Type == shared.StreamMessageFinished || msg.Type == shared.StreamMessageError || msg.Type == shared.StreamMessageAborted {
				body.Close()
				return
			}

		}
	}()
}

func readUntilSeparator(reader *bufio.Reader, separator string) (string, error) {
	var result []byte
	sepBytes := []byte(separator)
	for {
		b, err := reader.ReadByte()
		if err != nil {
			return string(result), err
		}
		result = append(result, b)
		if len(result) >= len(sepBytes) && bytes.HasSuffix(result, sepBytes) {
			return string(result[:len(result)-len(separator)]), nil
		}
	}
}