This repository has no description
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}