This repository has no description
1package executor
2
3import (
4 "context"
5 "crypto/rand"
6 "encoding/hex"
7 "encoding/json"
8 "fmt"
9 "os"
10 "time"
11
12 "google.golang.org/protobuf/proto"
13 "tangled.org/core/api/tangled"
14 "tangled.org/core/spindle/db"
15 millproto "tangled.org/core/spindle/mill/proto"
16 millv1 "tangled.org/core/spindle/mill/proto/gen"
17 "tangled.org/core/spindle/models"
18)
19
20const (
21 maxBatchBytes = 4 * 1024 * 1024
22 maxBatchEvents = 128
23)
24
25func generateEpoch() string {
26 var b [8]byte
27 _, _ = rand.Read(b[:])
28 return hex.EncodeToString(b[:])
29}
30
31func (e *Executor) initOutbox() error {
32 e.eventMu.Lock()
33 defer e.eventMu.Unlock()
34
35 epoch, _, err := e.db.GetOutboxState()
36 if err != nil {
37 return err
38 }
39 if epoch != "" {
40 e.epoch = epoch
41 } else {
42 e.epoch = generateEpoch()
43 if err := e.db.SetOutboxEpoch(e.epoch); err != nil {
44 return fmt.Errorf("set outbox epoch: %w", err)
45 }
46 }
47
48 rows, err := e.db.ListOutboxRows()
49 if err != nil {
50 return fmt.Errorf("list outbox rows: %w", err)
51 }
52
53 for _, r := range rows {
54 e.outboxBytes += r.ByteSize
55 }
56 _ = e.recoverPendingArtifacts()
57 return nil
58}
59
60func (e *Executor) appendAndSend(leaseID string, payload any, control bool) error {
61 entry := &millv1.Event{LeaseId: leaseID}
62 switch payload := payload.(type) {
63 case *millv1.Event_StatusEvent:
64 entry.Payload = payload
65 case *millv1.Event_AttemptResult:
66 entry.Payload = payload
67 default:
68 return fmt.Errorf("unsupported stream payload %T", payload)
69 }
70 isTerminal := false
71 if _, ok := payload.(*millv1.Event_AttemptResult); ok {
72 isTerminal = true
73 }
74 entry.Seqno = ^uint64(0)
75 wireSize := proto.Size(&millproto.Message{EventBatch: &millv1.EventBatch{
76 Epoch: e.epoch,
77 Events: []*millv1.Event{entry},
78 }})
79 entry.Seqno = 0
80 if wireSize > maxBatchBytes {
81 return fmt.Errorf("control stream entry exceeds wire limit: %d > %d", wireSize, maxBatchBytes)
82 }
83
84 encoded, err := proto.Marshal(entry)
85 if err != nil {
86 return fmt.Errorf("marshal stream entry: %w", err)
87 }
88
89 e.eventMu.Lock()
90 if control && !isTerminal && e.maxOutboxBytes > 0 && e.outboxBytes+int64(len(encoded)) > e.maxOutboxBytes {
91 e.l.Warn("outbox reserve exhausted; dropping nonterminal status", "cap", e.maxOutboxBytes)
92 e.eventMu.Unlock()
93 return nil
94 }
95 if _, err := e.db.AppendOutboxRow(encoded, control); err != nil {
96 e.eventMu.Unlock()
97 return fmt.Errorf("append outbox row: %w", err)
98 }
99 e.outboxBytes += int64(len(encoded))
100 e.eventMu.Unlock()
101
102 e.sendPending()
103 return nil
104}
105
106func (e *Executor) appendStatus(leaseID string, st *tangled.PipelineStatus) error {
107 if st.Status != string(models.StatusKindRunning) {
108 return fmt.Errorf("unsupported nonterminal status %q", st.Status)
109 }
110 exit, errStr := parseStatusExitAndError(st)
111 payload := &millv1.Event_StatusEvent{StatusEvent: &millv1.StatusEvent{
112 Status: millv1.NonterminalStatus_RUNNING,
113 Error: errStr,
114 ExitCode: exit,
115 }}
116 return e.appendAndSend(leaseID, payload, true)
117}
118
119func (e *Executor) appendTerminal(leaseID, status string, st *tangled.PipelineStatus) error {
120 return e.appendTerminalWithArtifact(leaseID, status, st, "", "")
121}
122
123func (e *Executor) appendTerminalWithArtifact(leaseID, status string, st *tangled.PipelineStatus, ref, hash string) error {
124 var terminalStatus millv1.TerminalStatus
125 switch status {
126 case string(models.StatusKindSuccess):
127 terminalStatus = millv1.TerminalStatus_SUCCESS
128 case string(models.StatusKindFailed):
129 terminalStatus = millv1.TerminalStatus_FAILED
130 case string(models.StatusKindTimeout):
131 terminalStatus = millv1.TerminalStatus_TIMEOUT
132 case string(models.StatusKindCancelled):
133 terminalStatus = millv1.TerminalStatus_CANCELLED
134 default:
135 return fmt.Errorf("unsupported terminal status %q", status)
136 }
137
138 exit, errStr := parseStatusExitAndError(st)
139 var logArtifact *millv1.LogArtifact
140 if ref != "" {
141 logArtifact = &millv1.LogArtifact{
142 Ref: ref,
143 Hash: hash,
144 }
145 }
146 payload := &millv1.Event_AttemptResult{AttemptResult: &millv1.AttemptResult{
147 Status: terminalStatus,
148 Error: errStr,
149 ExitCode: exit,
150 LogArtifact: logArtifact,
151 }}
152 return e.appendAndSend(leaseID, payload, true)
153}
154
155func (e *Executor) recoverPendingArtifacts() error {
156 if e.db == nil {
157 return nil
158 }
159 pending, err := e.db.ListPendingArtifacts()
160 if err != nil || len(pending) == 0 {
161 return err
162 }
163 for _, p := range pending {
164 if p.Ref != "" {
165 if e.writer == nil {
166 continue
167 }
168 ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
169 var logDir string
170 if e.cfg != nil {
171 logDir = e.cfg.Server.LogDir
172 }
173 logPath := models.LogFilePath(logDir, models.WorkflowId{Name: p.Workflow})
174 f, openErr := os.Open(logPath)
175 if openErr != nil {
176 cancel()
177 continue
178 }
179 uploadErr := e.writer.Put(ctx, p.Ref, f)
180 _ = f.Close()
181 cancel()
182 if uploadErr != nil {
183 continue
184 }
185 }
186
187 st := &tangled.PipelineStatus{
188 Status: p.Status,
189 Error: &p.Error,
190 ExitCode: &p.ExitCode,
191 }
192 if err := e.appendTerminalWithArtifact(p.LeaseID, p.Status, st, p.Ref, p.Hash); err == nil {
193 _ = e.db.RemovePendingArtifact(p.LeaseID)
194 }
195 }
196 return nil
197}
198
199func truncateEventString(value string) string {
200 if len(value) > 65536 {
201 return value[:65536]
202 }
203 return value
204}
205
206func (e *Executor) sendPending() {
207 e.connMu.Lock()
208 enc := e.enc
209 e.connMu.Unlock()
210
211 if enc == nil {
212 return
213 }
214
215 e.flushMu.Lock()
216 defer e.flushMu.Unlock()
217
218 if err := e.sendPendingLocked(enc); err != nil {
219 e.l.Error("send pending events failed", "err", err)
220 }
221}
222
223func (e *Executor) sendPendingLocked(enc messageEncoder) error {
224 for {
225 rows, err := e.db.ListOutboxRowsAfter(e.sentSeqno, maxBatchEvents)
226 if err != nil {
227 return fmt.Errorf("list outbox rows after %d: %w", e.sentSeqno, err)
228 }
229 if len(rows) == 0 {
230 return nil
231 }
232 if err := e.sendRows(enc, rows); err != nil {
233 return err
234 }
235 }
236}
237
238func (e *Executor) sendRows(enc messageEncoder, rows []db.OutboxRow) error {
239 events := make([]*millv1.Event, 0, len(rows))
240 var lastSeqno uint64
241
242 for _, r := range rows {
243 ev, err := decodeEvent(r.Payload, r.Seqno)
244 if err != nil {
245 return fmt.Errorf("decode event at seqno %d: %w", r.Seqno, err)
246 }
247 events = append(events, ev)
248 lastSeqno = r.Seqno
249 }
250
251 batch := &millv1.EventBatch{
252 Epoch: e.epoch,
253 Events: events,
254 }
255
256 e.sendMu.Lock()
257 defer e.sendMu.Unlock()
258
259 if err := enc.Encode(&millproto.Message{EventBatch: batch}); err != nil {
260 return fmt.Errorf("encode event batch (seqno %d..%d): %w", rows[0].Seqno, lastSeqno, err)
261 }
262 e.sentSeqno = lastSeqno
263 return nil
264}
265
266func (e *Executor) replay(ackSeqno uint64) error {
267 e.connMu.Lock()
268 enc := e.enc
269 e.connMu.Unlock()
270 if enc == nil {
271 return nil
272 }
273
274 e.flushMu.Lock()
275 defer e.flushMu.Unlock()
276
277 e.sendMu.Lock()
278 e.sentSeqno = ackSeqno
279 e.sendMu.Unlock()
280
281 return e.sendPendingLocked(enc)
282}
283
284func (e *Executor) handleAck(ack *millv1.Ack) {
285 if ack == nil || ack.GetEpoch() != e.epoch {
286 return
287 }
288
289 upTo := ack.GetUpToSeqno()
290 if upTo == 0 {
291 return
292 }
293
294 if err := e.deleteOutboxPrefix(upTo); err != nil {
295 e.l.Error("delete outbox prefix failed", "upTo", upTo, "err", err)
296 }
297}
298
299func (e *Executor) subtractOutboxBytes(deleted db.OutboxDeletion) {
300 e.outboxBytes = max(0, e.outboxBytes-deleted.Bytes)
301}
302
303func parseStatus(raw json.RawMessage) (*tangled.PipelineStatus, bool) {
304 var st tangled.PipelineStatus
305 if err := json.Unmarshal(raw, &st); err != nil {
306 return nil, false
307 }
308 return &st, true
309}
310
311func (e *Executor) deleteOutboxPrefix(upTo uint64) error {
312 e.eventMu.Lock()
313 defer e.eventMu.Unlock()
314
315 deleted, err := e.db.DeleteOutboxPrefix(upTo)
316 if err == nil {
317 e.subtractOutboxBytes(deleted)
318 }
319 return err
320}
321
322func parseStatusExitAndError(st *tangled.PipelineStatus) (int64, string) {
323 if st == nil {
324 return 0, ""
325 }
326 var exit int64
327 if st.ExitCode != nil {
328 exit = *st.ExitCode
329 }
330 var errStr string
331 if st.Error != nil {
332 errStr = truncateEventString(*st.Error)
333 }
334 return exit, errStr
335}
336
337func decodeEvent(payload []byte, seqno uint64) (*millv1.Event, error) {
338 var ev millv1.Event
339 if err := proto.Unmarshal(payload, &ev); err != nil {
340 return nil, err
341 }
342 ev.Seqno = seqno
343 return &ev, nil
344}