harness-drone/agent/updater.go

68 lines
1.7 KiB
Go
Raw Normal View History

package agent
import (
"fmt"
"github.com/Sirupsen/logrus"
"github.com/drone/drone/build"
2016-09-26 08:29:05 +00:00
"github.com/drone/drone/model"
2016-09-29 21:45:13 +00:00
"github.com/drone/mq/logger"
2016-09-26 08:29:05 +00:00
"github.com/drone/mq/stomp"
)
// UpdateFunc handles buid pipeline status updates.
2016-09-28 01:30:28 +00:00
type UpdateFunc func(*model.Work)
// LoggerFunc handles buid pipeline logging updates.
type LoggerFunc func(*build.Line)
2016-09-28 01:30:28 +00:00
var NoopUpdateFunc = func(*model.Work) {}
var TermLoggerFunc = func(line *build.Line) {
fmt.Println(line)
}
// NewClientUpdater returns an updater that sends updated build details
// to the drone server.
2016-09-26 08:29:05 +00:00
func NewClientUpdater(client *stomp.Client) UpdateFunc {
2016-09-28 01:30:28 +00:00
return func(w *model.Work) {
2016-09-26 08:29:05 +00:00
err := client.SendJSON("/queue/updates", w)
if err != nil {
2016-09-29 21:45:13 +00:00
logger.Warningf("Error updating %s/%s#%d.%d. %s",
w.Repo.Owner, w.Repo.Name, w.Build.Number, w.Job.Number, err)
}
2016-09-26 08:29:05 +00:00
if w.Job.Status != model.StatusRunning {
var dest = fmt.Sprintf("/topic/logs.%d", w.Job.ID)
var opts = []stomp.MessageOption{
stomp.WithHeader("eof", "true"),
stomp.WithRetain("all"),
}
2016-09-26 08:29:05 +00:00
if err := client.Send(dest, []byte("eof"), opts...); err != nil {
2016-09-29 21:45:13 +00:00
logger.Warningf("Error sending eof %s/%s#%d.%d. %s",
2016-09-26 08:29:05 +00:00
w.Repo.Owner, w.Repo.Name, w.Build.Number, w.Job.Number, err)
}
}
}
}
2016-09-26 08:29:05 +00:00
func NewClientLogger(client *stomp.Client, id int64, limit int64) LoggerFunc {
var size int64
2016-09-26 08:29:05 +00:00
var dest = fmt.Sprintf("/topic/logs.%d", id)
var opts = []stomp.MessageOption{
stomp.WithRetain("all"),
}
2016-09-26 08:29:05 +00:00
return func(line *build.Line) {
if size > limit {
return
}
2016-09-26 08:29:05 +00:00
if err := client.SendJSON(dest, line, opts...); err != nil {
logrus.Errorf("Error streaming build logs. %s", err)
}
size += int64(len(line.Out))
}
}