This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / spindle / mill / executor / outbox.go
8.1 kB 344 lines
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}