harness-drone/pkg/server/ws.go

136 lines
2.8 KiB
Go
Raw Normal View History

2015-04-08 22:43:59 +00:00
package server
import (
2015-05-19 04:53:34 +00:00
"fmt"
"io"
"net/http"
2015-04-30 02:57:43 +00:00
"strconv"
2015-04-08 22:43:59 +00:00
"github.com/drone/drone/pkg/bus"
2015-04-08 22:43:59 +00:00
2015-05-22 18:37:40 +00:00
log "github.com/drone/drone/Godeps/_workspace/src/github.com/Sirupsen/logrus"
"github.com/drone/drone/Godeps/_workspace/src/github.com/gin-gonic/gin"
"github.com/drone/drone/Godeps/_workspace/src/github.com/manucorporat/sse"
"github.com/drone/drone/pkg/docker"
2015-04-08 22:43:59 +00:00
)
// GetRepoEvents will upgrade the connection to a Websocket and will stream
// event updates to the browser.
func GetRepoEvents(c *gin.Context) {
bus_ := ToBus(c)
repo := ToRepo(c)
c.Writer.Header().Set("Content-Type", "text/event-stream")
2015-04-25 00:06:46 +00:00
eventc := make(chan *bus.Event, 1)
bus_.Subscribe(eventc)
2015-04-08 22:43:59 +00:00
defer func() {
bus_.Unsubscribe(eventc)
close(eventc)
log.Infof("closed event stream")
2015-04-08 22:43:59 +00:00
}()
c.Stream(func(w io.Writer) bool {
select {
case event := <-eventc:
if event == nil {
log.Infof("nil event received")
return false
}
if event.Kind == bus.EventRepo &&
event.Name == repo.FullName {
sse.Encode(w, sse.Event{
Event: "message",
Data: string(event.Msg),
})
}
case <-c.Writer.CloseNotify():
return false
}
return true
})
2015-04-08 22:43:59 +00:00
}
2015-04-30 02:57:43 +00:00
func GetStream(c *gin.Context) {
2015-05-19 04:53:34 +00:00
conf := ToSettings(c)
store := ToDatastore(c)
2015-04-30 02:57:43 +00:00
repo := ToRepo(c)
2015-05-06 07:56:06 +00:00
runner := ToRunner(c)
commitseq, _ := strconv.Atoi(c.Params.ByName("build"))
2015-06-19 00:36:52 +00:00
jobnum, _ := strconv.Atoi(c.Params.ByName("number"))
2015-04-30 02:57:43 +00:00
c.Writer.Header().Set("Content-Type", "text/event-stream")
2015-05-06 07:56:06 +00:00
build, err := store.BuildNumber(repo, commitseq)
if err != nil {
c.Fail(404, err)
return
}
job, err := store.JobNumber(build, jobnum)
if err != nil {
c.Fail(404, err)
return
}
var rc io.ReadCloser
// if the commit is being executed by an agent
// we'll proxy the build output directly to the
// remote Docker client, through the agent.
2015-05-28 07:18:46 +00:00
if conf.Agents.Secret != "" {
addr, err := store.Agent(build)
2015-05-19 04:53:34 +00:00
if err != nil {
c.Fail(500, err)
return
}
2015-06-19 00:36:52 +00:00
url := fmt.Sprintf("http://%s/stream/%d?token=%s", addr, job.ID, conf.Agents.Secret)
2015-05-19 04:53:34 +00:00
resp, err := http.Get(url)
if err != nil {
c.Fail(500, err)
return
} else if resp.StatusCode != 200 {
resp.Body.Close()
c.AbortWithStatus(resp.StatusCode)
return
}
rc = resp.Body
} else {
// else if the commit is not being executed
// by the build agent we can use the local runner
2015-06-19 00:36:52 +00:00
rc, err = runner.Logs(job)
if err != nil {
c.Fail(404, err)
return
}
2015-04-30 02:57:43 +00:00
}
defer func() {
rc.Close()
}()
2015-05-06 07:56:06 +00:00
go func() {
<-c.Writer.CloseNotify()
rc.Close()
2015-05-06 07:56:06 +00:00
}()
rw := &StreamWriter{c.Writer, 0}
2015-05-06 07:56:06 +00:00
docker.StdCopy(rw, rw, rc)
}
2015-05-06 07:56:06 +00:00
type StreamWriter struct {
writer gin.ResponseWriter
count int
2015-04-30 02:57:43 +00:00
}
func (w *StreamWriter) Write(data []byte) (int, error) {
var err = sse.Encode(w.writer, sse.Event{
Id: strconv.Itoa(w.count),
Event: "message",
Data: string(data),
2015-04-08 22:43:59 +00:00
})
w.writer.Flush()
w.count += len(data)
return len(data), err
2015-04-08 22:43:59 +00:00
}