2
0

operation.go 6.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371
  1. package redis_stream
  2. import (
  3. "context"
  4. "encoding/json"
  5. "fmt"
  6. "sync"
  7. "time"
  8. redis "github.com/go-redis/redis/v8"
  9. "github.com/porter-dev/porter/internal/models"
  10. "github.com/porter-dev/porter/provisioner/types"
  11. )
  12. func PushToOperationStream(
  13. client *redis.Client,
  14. infra *models.Infra,
  15. operation *models.Operation,
  16. data *types.TFResourceState,
  17. ) error {
  18. // pushes a state update to the state stream
  19. streamName := getStateStreamName(infra, operation)
  20. pushData := &types.TFResourceStateEntry{
  21. TFResourceState: data,
  22. PushedAt: time.Now(),
  23. }
  24. dataBytes, err := json.Marshal(pushData)
  25. if err != nil {
  26. return err
  27. }
  28. _, err = client.XAdd(context.TODO(), &redis.XAddArgs{
  29. Stream: streamName,
  30. ID: "*",
  31. Values: map[string]interface{}{
  32. "id": models.GetWorkspaceID(infra, operation),
  33. "data": dataBytes,
  34. },
  35. }).Result()
  36. return err
  37. }
  38. func SendOperationCompleted(
  39. client *redis.Client,
  40. infra *models.Infra,
  41. operation *models.Operation,
  42. ) error {
  43. // pushes a state update to the state stream
  44. streamName := getStateStreamName(infra, operation)
  45. data := map[string]interface{}{
  46. "status": "OPERATION_COMPLETED",
  47. }
  48. dataBytes, err := json.Marshal(data)
  49. if err != nil {
  50. return err
  51. }
  52. _, err = client.XAdd(context.TODO(), &redis.XAddArgs{
  53. Stream: streamName,
  54. ID: "*",
  55. Values: map[string]interface{}{
  56. "id": models.GetWorkspaceID(infra, operation),
  57. "data": dataBytes,
  58. },
  59. }).Result()
  60. return err
  61. }
  62. func PushToLogStream(
  63. client *redis.Client,
  64. infra *models.Infra,
  65. operation *models.Operation,
  66. data *types.TFLogLine,
  67. ) error {
  68. // pushes a state update to the state stream
  69. streamName := getLogsStreamName(infra, operation)
  70. _, err := client.XAdd(context.TODO(), &redis.XAddArgs{
  71. Stream: streamName,
  72. ID: "*",
  73. Values: map[string]interface{}{
  74. "log": getLogString(data),
  75. },
  76. }).Result()
  77. return err
  78. }
  79. func getLogString(data *types.TFLogLine) string {
  80. if data.Diagnostic.Detail != "" {
  81. return fmt.Sprintf("[%s] [%s] %s: %s\n", data.Level, data.Timestamp, data.Message, data.Diagnostic.Detail)
  82. }
  83. return fmt.Sprintf("[%s] [%s] %s\n", data.Level, data.Timestamp, data.Message)
  84. }
  85. type LogWriter func(log string) error
  86. func StreamOperationLogs(
  87. ctx context.Context,
  88. client *redis.Client,
  89. infra *models.Infra,
  90. operation *models.Operation,
  91. send LogWriter,
  92. ) error {
  93. streamName := getLogsStreamName(infra, operation)
  94. errorchan := make(chan error)
  95. redisCtx, cancel := context.WithCancel(context.Background())
  96. defer cancel()
  97. var wg sync.WaitGroup
  98. wg.Add(3)
  99. go func() {
  100. wg.Wait()
  101. close(errorchan)
  102. }()
  103. go func() {
  104. defer wg.Done()
  105. select {
  106. case <-ctx.Done():
  107. errorchan <- nil
  108. case <-redisCtx.Done():
  109. errorchan <- nil
  110. }
  111. }()
  112. go func() {
  113. defer wg.Done()
  114. // check intermittently that the stream still exists -- it may have been
  115. // cleaned up automatically
  116. failedCount := 0
  117. for {
  118. x, err := client.Exists(
  119. context.Background(),
  120. streamName,
  121. ).Result()
  122. // if the stream does not exist, increment the failed counter
  123. if x == 0 || err != nil {
  124. failedCount++
  125. } else {
  126. failedCount = 0
  127. }
  128. if failedCount >= 2 {
  129. errorchan <- nil
  130. return
  131. }
  132. // wait 5 seconds in between pings
  133. time.Sleep(5 * time.Second)
  134. }
  135. }()
  136. go func() {
  137. defer wg.Done()
  138. lastID := "0-0"
  139. for {
  140. if redisCtx.Err() != nil {
  141. errorchan <- nil
  142. return
  143. }
  144. xstream, err := client.XRead(
  145. redisCtx,
  146. &redis.XReadArgs{
  147. Streams: []string{streamName, lastID},
  148. Block: 0,
  149. },
  150. ).Result()
  151. if err != nil {
  152. errorchan <- err
  153. return
  154. }
  155. messages := xstream[0].Messages
  156. lastID = messages[len(messages)-1].ID
  157. for _, msg := range messages {
  158. dataInter, ok := msg.Values["log"]
  159. if !ok {
  160. continue
  161. }
  162. dataString, ok := dataInter.(string)
  163. if !ok {
  164. continue
  165. }
  166. err = send(dataString)
  167. if err != nil {
  168. errorchan <- err
  169. return
  170. }
  171. }
  172. }
  173. }()
  174. var err error
  175. for err = range errorchan {
  176. cancel()
  177. }
  178. return err
  179. }
  180. type StateUpdateWriter func(update *types.TFResourceState) error
  181. func StreamStateUpdate(
  182. ctx context.Context,
  183. client *redis.Client,
  184. infra *models.Infra,
  185. operation *models.Operation,
  186. send StateUpdateWriter,
  187. ) error {
  188. streamName := getStateStreamName(infra, operation)
  189. errorchan := make(chan error)
  190. redisCtx, cancel := context.WithCancel(context.Background())
  191. defer cancel()
  192. var wg sync.WaitGroup
  193. wg.Add(3)
  194. go func() {
  195. wg.Wait()
  196. close(errorchan)
  197. }()
  198. go func() {
  199. defer wg.Done()
  200. select {
  201. case <-ctx.Done():
  202. errorchan <- nil
  203. case <-redisCtx.Done():
  204. errorchan <- nil
  205. }
  206. }()
  207. go func() {
  208. defer wg.Done()
  209. // check intermittently that the stream still exists -- it may have been
  210. // cleaned up automatically
  211. failedCount := 0
  212. for {
  213. x, err := client.Exists(
  214. context.Background(),
  215. streamName,
  216. ).Result()
  217. // if the stream does not exist, increment the failed counter
  218. if x == 0 || err != nil {
  219. failedCount++
  220. } else {
  221. failedCount = 0
  222. }
  223. if failedCount >= 2 {
  224. errorchan <- nil
  225. return
  226. }
  227. // wait 5 seconds in between pings
  228. time.Sleep(5 * time.Second)
  229. }
  230. }()
  231. go func() {
  232. defer wg.Done()
  233. lastID := "0-0"
  234. for {
  235. if redisCtx.Err() != nil {
  236. errorchan <- nil
  237. return
  238. }
  239. xstream, err := client.XRead(
  240. ctx,
  241. &redis.XReadArgs{
  242. Streams: []string{streamName, lastID},
  243. Block: 0,
  244. },
  245. ).Result()
  246. if err != nil {
  247. errorchan <- err
  248. return
  249. }
  250. messages := xstream[0].Messages
  251. lastID = messages[len(messages)-1].ID
  252. for _, msg := range messages {
  253. stateData := &types.TFResourceState{}
  254. dataInter, ok := msg.Values["data"]
  255. if !ok {
  256. continue
  257. }
  258. dataString, ok := dataInter.(string)
  259. if !ok {
  260. continue
  261. }
  262. err := json.Unmarshal([]byte(dataString), stateData)
  263. if err != nil {
  264. continue
  265. }
  266. err = send(stateData)
  267. if err != nil {
  268. errorchan <- err
  269. return
  270. }
  271. }
  272. }
  273. }()
  274. var err error
  275. for err = range errorchan {
  276. cancel()
  277. }
  278. return err
  279. }
  280. func getStateStreamName(
  281. infra *models.Infra,
  282. operation *models.Operation,
  283. ) string {
  284. return fmt.Sprintf("%s-state", models.GetWorkspaceID(infra, operation))
  285. }
  286. func getLogsStreamName(
  287. infra *models.Infra,
  288. operation *models.Operation,
  289. ) string {
  290. return fmt.Sprintf("%s-logs", models.GetWorkspaceID(infra, operation))
  291. }
  292. func getLogsFileName(
  293. infra *models.Infra,
  294. operation *models.Operation,
  295. ) string {
  296. return fmt.Sprintf("%s-logs.txt", models.GetWorkspaceID(infra, operation))
  297. }