(task TaskRunner)
| 98 | func (m *TaskBase) CloseFinal() error { return nil } |
| 99 | |
| 100 | func MakeHandler(task TaskRunner) MessageHandler { |
| 101 | out := task.MessageOut() |
| 102 | return func(ctx *plan.Context, msg schema.Message) bool { |
| 103 | select { |
| 104 | case out <- msg: |
| 105 | return true |
| 106 | case <-task.SigChan(): |
| 107 | return false |
| 108 | } |
| 109 | } |
| 110 | } |
| 111 | |
| 112 | func (m *TaskBase) Run() error { |
| 113 | defer m.Ctx.Recover() // Our context can recover panics, save error msg |
nothing calls this directly
no test coverage detected