agola/internal/services/runservice/api/api.go

760 lines
18 KiB
Go
Raw Normal View History

2019-02-21 14:54:50 +00:00
// Copyright 2019 Sorint.lab
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied
// See the License for the specific language governing permissions and
// limitations under the License.
package api
import (
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"strconv"
2019-07-01 09:40:20 +00:00
"agola.io/agola/internal/datamanager"
"agola.io/agola/internal/db"
"agola.io/agola/internal/etcd"
"agola.io/agola/internal/objectstorage"
"agola.io/agola/internal/services/runservice/action"
"agola.io/agola/internal/services/runservice/common"
"agola.io/agola/internal/services/runservice/readdb"
"agola.io/agola/internal/services/runservice/store"
"agola.io/agola/internal/util"
rsapitypes "agola.io/agola/services/runservice/api/types"
"agola.io/agola/services/runservice/types"
2019-02-21 14:54:50 +00:00
"github.com/gorilla/mux"
2019-05-15 07:46:21 +00:00
etcdclientv3 "go.etcd.io/etcd/clientv3"
etcdclientv3rpc "go.etcd.io/etcd/etcdserver/api/v3rpc/rpctypes"
"go.etcd.io/etcd/mvcc/mvccpb"
2019-02-21 14:54:50 +00:00
"go.uber.org/zap"
errors "golang.org/x/xerrors"
2019-02-21 14:54:50 +00:00
)
type ErrorResponse struct {
Message string `json:"message"`
}
func ErrorResponseFromError(err error) *ErrorResponse {
var aerr error
// use inner errors if of these types
switch {
case util.IsBadRequest(err):
var cerr *util.ErrBadRequest
errors.As(err, &cerr)
aerr = cerr
case util.IsNotExist(err):
var cerr *util.ErrNotExist
errors.As(err, &cerr)
aerr = cerr
case util.IsForbidden(err):
var cerr *util.ErrForbidden
errors.As(err, &cerr)
aerr = cerr
case util.IsUnauthorized(err):
var cerr *util.ErrUnauthorized
errors.As(err, &cerr)
aerr = cerr
case util.IsInternal(err):
var cerr *util.ErrInternal
errors.As(err, &cerr)
aerr = cerr
}
if aerr != nil {
return &ErrorResponse{Message: aerr.Error()}
}
// on generic error return an generic message to not leak the real error
return &ErrorResponse{Message: "internal server error"}
}
func httpError(w http.ResponseWriter, err error) bool {
if err == nil {
return false
}
response := ErrorResponseFromError(err)
resj, merr := json.Marshal(response)
if merr != nil {
w.WriteHeader(http.StatusInternalServerError)
return true
}
switch {
case util.IsBadRequest(err):
w.WriteHeader(http.StatusBadRequest)
_, _ = w.Write(resj)
case util.IsNotExist(err):
w.WriteHeader(http.StatusNotFound)
_, _ = w.Write(resj)
case util.IsForbidden(err):
w.WriteHeader(http.StatusForbidden)
_, _ = w.Write(resj)
case util.IsUnauthorized(err):
w.WriteHeader(http.StatusUnauthorized)
_, _ = w.Write(resj)
case util.IsInternal(err):
w.WriteHeader(http.StatusInternalServerError)
_, _ = w.Write(resj)
default:
w.WriteHeader(http.StatusInternalServerError)
_, _ = w.Write(resj)
}
return true
}
func httpResponse(w http.ResponseWriter, code int, res interface{}) error {
w.Header().Set("Content-Type", "application/json")
if res != nil {
resj, err := json.Marshal(res)
if err != nil {
httpError(w, err)
return err
}
w.WriteHeader(code)
_, err = w.Write(resj)
return err
}
w.WriteHeader(code)
return nil
}
2019-02-21 14:54:50 +00:00
type LogsHandler struct {
log *zap.SugaredLogger
e *etcd.Store
ost *objectstorage.ObjStorage
dm *datamanager.DataManager
2019-02-21 14:54:50 +00:00
}
func NewLogsHandler(logger *zap.Logger, e *etcd.Store, ost *objectstorage.ObjStorage, dm *datamanager.DataManager) *LogsHandler {
2019-02-21 14:54:50 +00:00
return &LogsHandler{
log: logger.Sugar(),
e: e,
ost: ost,
dm: dm,
2019-02-21 14:54:50 +00:00
}
}
func (h *LogsHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
ctx := r.Context()
2019-03-13 14:48:35 +00:00
q := r.URL.Query()
2019-02-21 14:54:50 +00:00
2019-03-13 14:48:35 +00:00
runID := q.Get("runid")
2019-02-21 14:54:50 +00:00
if runID == "" {
http.Error(w, "", http.StatusBadRequest)
return
}
2019-03-13 14:48:35 +00:00
taskID := q.Get("taskid")
2019-02-21 14:54:50 +00:00
if taskID == "" {
http.Error(w, "", http.StatusBadRequest)
return
}
2019-03-13 14:48:35 +00:00
_, setup := q["setup"]
stepStr := q.Get("step")
if !setup && stepStr == "" {
2019-02-21 14:54:50 +00:00
http.Error(w, "", http.StatusBadRequest)
return
}
2019-03-13 14:48:35 +00:00
if setup && stepStr != "" {
2019-02-21 14:54:50 +00:00
http.Error(w, "", http.StatusBadRequest)
return
}
2019-03-13 14:48:35 +00:00
var step int
if stepStr != "" {
var err error
step, err = strconv.Atoi(stepStr)
if err != nil {
http.Error(w, "", http.StatusBadRequest)
return
}
}
2019-02-21 14:54:50 +00:00
follow := false
2019-03-13 14:48:35 +00:00
if _, ok := q["follow"]; ok {
2019-02-21 14:54:50 +00:00
follow = true
}
if err, sendError := h.readTaskLogs(ctx, runID, taskID, setup, step, w, follow); err != nil {
2019-02-21 14:54:50 +00:00
h.log.Errorf("err: %+v", err)
if sendError {
switch {
case util.IsNotExist(err):
httpError(w, util.NewErrNotExist(errors.Errorf("log doesn't exist: %w", err)))
2019-02-21 14:54:50 +00:00
default:
httpError(w, err)
2019-02-21 14:54:50 +00:00
}
}
}
}
func (h *LogsHandler) readTaskLogs(ctx context.Context, runID, taskID string, setup bool, step int, w http.ResponseWriter, follow bool) (error, bool) {
r, err := store.GetRunEtcdOrOST(ctx, h.e, h.dm, runID)
2019-02-21 14:54:50 +00:00
if err != nil {
return err, true
}
if r == nil {
return util.NewErrNotExist(errors.Errorf("no such run with id: %s", runID)), true
2019-02-21 14:54:50 +00:00
}
task, ok := r.Tasks[taskID]
2019-02-21 14:54:50 +00:00
if !ok {
return util.NewErrNotExist(errors.Errorf("no such task with ID %s in run %s", taskID, runID)), true
2019-02-21 14:54:50 +00:00
}
if len(task.Steps) <= step {
return util.NewErrNotExist(errors.Errorf("no such step for task %s in run %s", taskID, runID)), true
2019-02-21 14:54:50 +00:00
}
// if the log has been already fetched use it, otherwise fetch it from the executor
if task.Steps[step].LogPhase == types.RunTaskFetchPhaseFinished {
2019-03-13 14:48:35 +00:00
var logPath string
if setup {
logPath = store.OSTRunTaskSetupLogPath(task.ID)
2019-03-13 14:48:35 +00:00
} else {
logPath = store.OSTRunTaskStepLogPath(task.ID, step)
2019-03-13 14:48:35 +00:00
}
f, err := h.ost.ReadObject(logPath)
2019-02-21 14:54:50 +00:00
if err != nil {
if objectstorage.IsNotExist(err) {
return util.NewErrNotExist(err), true
2019-02-21 14:54:50 +00:00
}
return err, true
}
defer f.Close()
return sendLogs(w, f), false
2019-02-21 14:54:50 +00:00
}
et, err := store.GetExecutorTask(ctx, h.e, task.ID)
if err != nil {
if err == etcd.ErrKeyNotFound {
return util.NewErrNotExist(errors.Errorf("executor task with id %q doesn't exist", task.ID)), true
}
2019-02-21 14:54:50 +00:00
return err, true
}
executor, err := store.GetExecutor(ctx, h.e, et.Spec.ExecutorID)
if err != nil {
if err == etcd.ErrKeyNotFound {
return util.NewErrNotExist(errors.Errorf("executor with id %q doesn't exist", et.Spec.ExecutorID)), true
}
2019-02-21 14:54:50 +00:00
return err, true
}
2019-03-13 14:48:35 +00:00
var url string
if setup {
url = fmt.Sprintf("%s/api/v1alpha/executor/logs?taskid=%s&setup", executor.ListenURL, taskID)
} else {
url = fmt.Sprintf("%s/api/v1alpha/executor/logs?taskid=%s&step=%d", executor.ListenURL, taskID, step)
}
2019-02-21 14:54:50 +00:00
if follow {
url += "&follow"
}
req, err := http.Get(url)
if err != nil {
return err, true
}
defer req.Body.Close()
if req.StatusCode != http.StatusOK {
if req.StatusCode == http.StatusNotFound {
return util.NewErrNotExist(errors.New("no log on executor")), true
2019-02-21 14:54:50 +00:00
}
return errors.Errorf("received http status: %d", req.StatusCode), true
}
// write and flush the headers so the client will receive the response
// header also if there're currently no lines to send
w.Header().Set("Cache-Control", "no-cache")
w.Header().Set("Connection", "keep-alive")
w.WriteHeader(http.StatusOK)
var flusher http.Flusher
if fl, ok := w.(http.Flusher); ok {
flusher = fl
}
if flusher != nil {
flusher.Flush()
}
return sendLogs(w, req.Body), false
2019-02-21 14:54:50 +00:00
}
func sendLogs(w http.ResponseWriter, r io.Reader) error {
buf := make([]byte, 406)
2019-02-21 14:54:50 +00:00
var flusher http.Flusher
if fl, ok := w.(http.Flusher); ok {
flusher = fl
}
stop := false
for {
if stop {
return nil
}
n, err := r.Read(buf)
//data, err := br.ReadBytes('\n')
2019-02-21 14:54:50 +00:00
if err != nil {
if err != io.EOF {
return err
}
if n == 0 {
2019-02-21 14:54:50 +00:00
return nil
}
stop = true
}
if _, err := w.Write(buf[:n]); err != nil {
return err
2019-02-21 14:54:50 +00:00
}
if flusher != nil {
flusher.Flush()
}
}
}
type ChangeGroupsUpdateTokensHandler struct {
log *zap.SugaredLogger
readDB *readdb.ReadDB
}
func NewChangeGroupsUpdateTokensHandler(logger *zap.Logger, readDB *readdb.ReadDB) *ChangeGroupsUpdateTokensHandler {
return &ChangeGroupsUpdateTokensHandler{
log: logger.Sugar(),
readDB: readDB,
}
}
func (h *ChangeGroupsUpdateTokensHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
ctx := r.Context()
2019-02-21 14:54:50 +00:00
query := r.URL.Query()
groups := query["changegroup"]
var cgt *types.ChangeGroupsUpdateToken
err := h.readDB.Do(ctx, func(tx *db.Tx) error {
2019-02-21 14:54:50 +00:00
var err error
cgt, err = h.readDB.GetChangeGroupsUpdateTokens(tx, groups)
return err
})
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
cgts, err := types.MarshalChangeGroupsUpdateToken(cgt)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
if err := httpResponse(w, http.StatusOK, cgts); err != nil {
h.log.Errorf("err: %+v", err)
2019-02-21 14:54:50 +00:00
}
}
type RunHandler struct {
log *zap.SugaredLogger
e *etcd.Store
dm *datamanager.DataManager
2019-02-21 14:54:50 +00:00
readDB *readdb.ReadDB
}
func NewRunHandler(logger *zap.Logger, e *etcd.Store, dm *datamanager.DataManager, readDB *readdb.ReadDB) *RunHandler {
2019-02-21 14:54:50 +00:00
return &RunHandler{
log: logger.Sugar(),
e: e,
dm: dm,
2019-02-21 14:54:50 +00:00
readDB: readDB,
}
}
func (h *RunHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
ctx := r.Context()
2019-02-21 14:54:50 +00:00
vars := mux.Vars(r)
runID := vars["runid"]
2019-05-06 13:18:49 +00:00
query := r.URL.Query()
changeGroups := query["changegroup"]
var run *types.Run
var cgt *types.ChangeGroupsUpdateToken
err := h.readDB.Do(ctx, func(tx *db.Tx) error {
2019-05-06 13:18:49 +00:00
var err error
run, err = h.readDB.GetRun(tx, runID)
if err != nil {
h.log.Errorf("err: %+v", err)
return err
}
cgt, err = h.readDB.GetChangeGroupsUpdateTokens(tx, changeGroups)
return err
})
if err != nil {
2019-02-21 14:54:50 +00:00
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
if run == nil {
httpError(w, util.NewErrNotExist(errors.Errorf("run %q doesn't exist", runID)))
return
}
2019-05-06 13:18:49 +00:00
cgts, err := types.MarshalChangeGroupsUpdateToken(cgt)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
2019-02-21 14:54:50 +00:00
return
}
rc, err := store.OSTGetRunConfig(h.dm, run.ID)
2019-02-21 14:54:50 +00:00
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
res := &rsapitypes.RunResponse{
2019-05-06 13:18:49 +00:00
Run: run,
RunConfig: rc,
ChangeGroupsUpdateToken: cgts,
2019-02-21 14:54:50 +00:00
}
if err := httpResponse(w, http.StatusOK, res); err != nil {
h.log.Errorf("err: %+v", err)
2019-02-21 14:54:50 +00:00
}
}
const (
DefaultRunsLimit = 25
MaxRunsLimit = 40
)
type RunsHandler struct {
log *zap.SugaredLogger
readDB *readdb.ReadDB
}
func NewRunsHandler(logger *zap.Logger, readDB *readdb.ReadDB) *RunsHandler {
return &RunsHandler{
log: logger.Sugar(),
readDB: readDB,
}
}
func (h *RunsHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
ctx := r.Context()
2019-02-21 14:54:50 +00:00
query := r.URL.Query()
phaseFilter := types.RunPhaseFromStringSlice(query["phase"])
resultFilter := types.RunResultFromStringSlice(query["result"])
2019-02-21 14:54:50 +00:00
changeGroups := query["changegroup"]
groups := query["group"]
2019-03-29 11:00:18 +00:00
_, lastRun := query["lastrun"]
2019-02-21 14:54:50 +00:00
limitS := query.Get("limit")
limit := DefaultRunsLimit
if limitS != "" {
var err error
limit, err = strconv.Atoi(limitS)
if err != nil {
http.Error(w, "", http.StatusBadRequest)
return
}
}
if limit < 0 {
http.Error(w, "limit must be greater or equal than 0", http.StatusBadRequest)
return
}
if limit > MaxRunsLimit {
limit = MaxRunsLimit
}
sortOrder := types.SortOrderDesc
if _, ok := query["asc"]; ok {
sortOrder = types.SortOrderAsc
}
start := query.Get("start")
var runs []*types.Run
var cgt *types.ChangeGroupsUpdateToken
err := h.readDB.Do(ctx, func(tx *db.Tx) error {
2019-02-21 14:54:50 +00:00
var err error
runs, err = h.readDB.GetRuns(tx, groups, lastRun, phaseFilter, resultFilter, start, limit, sortOrder)
2019-02-21 14:54:50 +00:00
if err != nil {
h.log.Errorf("err: %+v", err)
return err
}
cgt, err = h.readDB.GetChangeGroupsUpdateTokens(tx, changeGroups)
return err
})
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
cgts, err := types.MarshalChangeGroupsUpdateToken(cgt)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
res := &rsapitypes.GetRunsResponse{
2019-02-21 14:54:50 +00:00
Runs: runs,
ChangeGroupsUpdateToken: cgts,
}
if err := httpResponse(w, http.StatusOK, res); err != nil {
h.log.Errorf("err: %+v", err)
2019-02-21 14:54:50 +00:00
}
}
type RunCreateHandler struct {
log *zap.SugaredLogger
ah *action.ActionHandler
2019-02-21 14:54:50 +00:00
}
func NewRunCreateHandler(logger *zap.Logger, ah *action.ActionHandler) *RunCreateHandler {
2019-02-21 14:54:50 +00:00
return &RunCreateHandler{
log: logger.Sugar(),
ah: ah,
2019-02-21 14:54:50 +00:00
}
}
func (h *RunCreateHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
ctx := r.Context()
var req rsapitypes.RunCreateRequest
2019-02-21 14:54:50 +00:00
d := json.NewDecoder(r.Body)
if err := d.Decode(&req); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
creq := &action.RunCreateRequest{
RunConfigTasks: req.RunConfigTasks,
Name: req.Name,
Group: req.Group,
SetupErrors: req.SetupErrors,
StaticEnvironment: req.StaticEnvironment,
CacheGroup: req.CacheGroup,
RunID: req.RunID,
FromStart: req.FromStart,
ResetTasks: req.ResetTasks,
2019-02-21 14:54:50 +00:00
Environment: req.Environment,
Annotations: req.Annotations,
ChangeGroupsUpdateToken: req.ChangeGroupsUpdateToken,
}
rb, err := h.ah.CreateRun(ctx, creq)
if err != nil {
h.log.Errorf("err: %+v", err)
httpError(w, err)
2019-02-21 14:54:50 +00:00
return
}
res := &rsapitypes.RunResponse{
Run: rb.Run,
RunConfig: rb.Rc,
}
if err := httpResponse(w, http.StatusCreated, res); err != nil {
h.log.Errorf("err: %+v", err)
}
2019-02-21 14:54:50 +00:00
}
type RunActionsHandler struct {
log *zap.SugaredLogger
ah *action.ActionHandler
2019-02-21 14:54:50 +00:00
}
func NewRunActionsHandler(logger *zap.Logger, ah *action.ActionHandler) *RunActionsHandler {
2019-02-21 14:54:50 +00:00
return &RunActionsHandler{
log: logger.Sugar(),
ah: ah,
2019-02-21 14:54:50 +00:00
}
}
func (h *RunActionsHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
ctx := r.Context()
vars := mux.Vars(r)
runID := vars["runid"]
var req rsapitypes.RunActionsRequest
2019-02-21 14:54:50 +00:00
d := json.NewDecoder(r.Body)
if err := d.Decode(&req); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
switch req.ActionType {
case rsapitypes.RunActionTypeChangePhase:
creq := &action.RunChangePhaseRequest{
2019-02-21 14:54:50 +00:00
RunID: runID,
Phase: req.Phase,
ChangeGroupsUpdateToken: req.ChangeGroupsUpdateToken,
}
if err := h.ah.ChangeRunPhase(ctx, creq); err != nil {
h.log.Errorf("err: %+v", err)
httpError(w, err)
2019-02-21 14:54:50 +00:00
return
}
case rsapitypes.RunActionTypeStop:
creq := &action.RunStopRequest{
2019-03-08 09:02:37 +00:00
RunID: runID,
ChangeGroupsUpdateToken: req.ChangeGroupsUpdateToken,
}
if err := h.ah.StopRun(ctx, creq); err != nil {
h.log.Errorf("err: %+v", err)
httpError(w, err)
2019-03-08 09:02:37 +00:00
return
}
2019-02-21 14:54:50 +00:00
default:
http.Error(w, "", http.StatusBadRequest)
return
}
}
type RunTaskActionsHandler struct {
log *zap.SugaredLogger
ah *action.ActionHandler
2019-02-21 14:54:50 +00:00
}
func NewRunTaskActionsHandler(logger *zap.Logger, ah *action.ActionHandler) *RunTaskActionsHandler {
2019-02-21 14:54:50 +00:00
return &RunTaskActionsHandler{
log: logger.Sugar(),
ah: ah,
2019-02-21 14:54:50 +00:00
}
}
func (h *RunTaskActionsHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
ctx := r.Context()
vars := mux.Vars(r)
runID := vars["runid"]
taskID := vars["taskid"]
var req rsapitypes.RunTaskActionsRequest
2019-02-21 14:54:50 +00:00
d := json.NewDecoder(r.Body)
if err := d.Decode(&req); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
switch req.ActionType {
case rsapitypes.RunTaskActionTypeSetAnnotations:
creq := &action.RunTaskSetAnnotationsRequest{
RunID: runID,
TaskID: taskID,
Annotations: req.Annotations,
ChangeGroupsUpdateToken: req.ChangeGroupsUpdateToken,
}
if err := h.ah.RunTaskSetAnnotations(ctx, creq); err != nil {
h.log.Errorf("err: %+v", err)
httpError(w, err)
return
}
case rsapitypes.RunTaskActionTypeApprove:
creq := &action.RunTaskApproveRequest{
2019-02-21 14:54:50 +00:00
RunID: runID,
TaskID: taskID,
ChangeGroupsUpdateToken: req.ChangeGroupsUpdateToken,
}
if err := h.ah.ApproveRunTask(ctx, creq); err != nil {
h.log.Errorf("err: %+v", err)
httpError(w, err)
2019-02-21 14:54:50 +00:00
return
}
2019-02-21 14:54:50 +00:00
default:
http.Error(w, "", http.StatusBadRequest)
return
}
}
2019-05-15 07:46:21 +00:00
type RunEventsHandler struct {
log *zap.SugaredLogger
e *etcd.Store
ost *objectstorage.ObjStorage
dm *datamanager.DataManager
}
func NewRunEventsHandler(logger *zap.Logger, e *etcd.Store, ost *objectstorage.ObjStorage, dm *datamanager.DataManager) *RunEventsHandler {
return &RunEventsHandler{
log: logger.Sugar(),
e: e,
ost: ost,
dm: dm,
}
}
//
func (h *RunEventsHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
ctx := r.Context()
q := r.URL.Query()
// TODO(sgotti) handle additional events filtering (by type, etc...)
startRunEventID := q.Get("startruneventid")
if err := h.sendRunEvents(ctx, startRunEventID, w); err != nil {
h.log.Errorf("err: %+v", err)
}
}
func (h *RunEventsHandler) sendRunEvents(ctx context.Context, startRunEventID string, w http.ResponseWriter) error {
w.Header().Set("Content-Type", "text/event-stream")
w.Header().Set("Cache-Control", "no-cache")
w.Header().Set("Connection", "keep-alive")
var flusher http.Flusher
if fl, ok := w.(http.Flusher); ok {
flusher = fl
}
// TODO(sgotti) fetch from previous events (handle startRunEventID).
// Use the readdb instead of etcd
wctx, cancel := context.WithCancel(ctx)
defer cancel()
wctx = etcdclientv3.WithRequireLeader(wctx)
wch := h.e.WatchKey(wctx, common.EtcdRunEventKey, 0)
for wresp := range wch {
if wresp.Canceled {
err := wresp.Err()
if err == etcdclientv3rpc.ErrCompacted {
h.log.Errorf("required events already compacted")
}
return errors.Errorf("watch error: %w", err)
2019-05-15 07:46:21 +00:00
}
for _, ev := range wresp.Events {
switch ev.Type {
case mvccpb.PUT:
var runEvent *types.RunEvent
if err := json.Unmarshal(ev.Kv.Value, &runEvent); err != nil {
return errors.Errorf("failed to unmarshal run: %w", err)
2019-05-15 07:46:21 +00:00
}
if _, err := w.Write([]byte(fmt.Sprintf("data: %s\n\n", ev.Kv.Value))); err != nil {
return err
}
}
}
if flusher != nil {
flusher.Flush()
}
}
return nil
}