This repository has no description
0

Configure Feed

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

core / spindle / db / events.go
3.2 kB 123 lines
1package db 2 3import ( 4 "encoding/json" 5 "time" 6 7 "tangled.org/core/api/tangled" 8 "tangled.org/core/eventstream" 9 "tangled.org/core/notifier" 10 "tangled.org/core/spindle/models" 11 "tangled.org/core/tid" 12) 13 14func (d *DB) insertEvent(event eventstream.Event, n *notifier.Notifier) error { 15 return eventstream.Insert(d, event, n) 16} 17 18func (d *DB) GetEvents(cursor int64, limit int) ([]eventstream.Event, error) { 19 return eventstream.List(d, cursor, limit) 20} 21 22func (d *DB) CreatePipelineEvent(rkey string, pipeline tangled.Pipeline, n *notifier.Notifier) error { 23 eventJson, err := json.Marshal(pipeline) 24 if err != nil { 25 return err 26 } 27 event := eventstream.Event{ 28 Rkey: rkey, 29 Nsid: tangled.PipelineNSID, 30 EventJson: eventJson, 31 } 32 return d.insertEvent(event, n) 33} 34 35func (d *DB) createStatusEvent( 36 workflowId models.WorkflowId, 37 statusKind models.StatusKind, 38 workflowError *string, 39 exitCode *int64, 40 n *notifier.Notifier, 41) error { 42 now := time.Now() 43 pipelineAtUri := workflowId.PipelineId.AtUri() 44 s := tangled.PipelineStatus{ 45 CreatedAt: now.Format(time.RFC3339), 46 Error: workflowError, 47 ExitCode: exitCode, 48 Pipeline: string(pipelineAtUri), 49 Workflow: workflowId.Name, 50 Status: string(statusKind), 51 } 52 53 eventJson, err := json.Marshal(s) 54 if err != nil { 55 return err 56 } 57 58 event := eventstream.Event{ 59 Rkey: tid.TID(), 60 Nsid: tangled.PipelineStatusNSID, 61 EventJson: eventJson, 62 } 63 64 return d.insertEvent(event, n) 65} 66 67func (d *DB) GetStatus(workflowId models.WorkflowId) (*tangled.PipelineStatus, error) { 68 pipelineAtUri := workflowId.PipelineId.AtUri() 69 70 var eventJson string 71 err := d.QueryRow( 72 ` 73 select 74 event from events 75 where 76 nsid = ? 77 and json_extract(event, '$.pipeline') = ? 78 and json_extract(event, '$.workflow') = ? 79 order by 80 created desc 81 limit 82 1 83 `, 84 tangled.PipelineStatusNSID, 85 string(pipelineAtUri), 86 workflowId.Name, 87 ).Scan(&eventJson) 88 89 if err != nil { 90 return nil, err 91 } 92 93 var status tangled.PipelineStatus 94 if err := json.Unmarshal([]byte(eventJson), &status); err != nil { 95 return nil, err 96 } 97 98 return &status, nil 99} 100 101func (d *DB) StatusPending(workflowId models.WorkflowId, n *notifier.Notifier) error { 102 return d.createStatusEvent(workflowId, models.StatusKindPending, nil, nil, n) 103} 104 105func (d *DB) StatusRunning(workflowId models.WorkflowId, n *notifier.Notifier) error { 106 return d.createStatusEvent(workflowId, models.StatusKindRunning, nil, nil, n) 107} 108 109func (d *DB) StatusFailed(workflowId models.WorkflowId, workflowError string, exitCode int64, n *notifier.Notifier) error { 110 return d.createStatusEvent(workflowId, models.StatusKindFailed, &workflowError, &exitCode, n) 111} 112 113func (d *DB) StatusCancelled(workflowId models.WorkflowId, workflowError string, exitCode int64, n *notifier.Notifier) error { 114 return d.createStatusEvent(workflowId, models.StatusKindCancelled, &workflowError, &exitCode, n) 115} 116 117func (d *DB) StatusSuccess(workflowId models.WorkflowId, n *notifier.Notifier) error { 118 return d.createStatusEvent(workflowId, models.StatusKindSuccess, nil, nil, n) 119} 120 121func (d *DB) StatusTimeout(workflowId models.WorkflowId, n *notifier.Notifier) error { 122 return d.createStatusEvent(workflowId, models.StatusKindTimeout, nil, nil, n) 123}