Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330package db
import ( "context" "database/sql" "encoding/json" "fmt" "tangled.org/core/api/tangled" "tangled.org/core/eventstream" "tangled.org/core/notifier" "tangled.org/core/spindle/models" "tangled.org/core/tid" "time")
func (d *DB) insertEvent(event eventstream.Event, n *notifier.Notifier) error { return eventstream.Insert(d, event, n)}
func (d *DB) GetEvents(cursor int64, limit int) ([]eventstream.Event, error) { return eventstream.List(d, cursor, limit)}
func (d *DB) EventHighWater() (int64, error) { return eventstream.HighWater(d)}
func (d *DB) CreatePipelineEvent(rkey string, pipeline tangled.Pipeline, n *notifier.Notifier) error { eventJson, err := json.Marshal(pipeline) if err != nil { return err } event := eventstream.Event{ Rkey: rkey, Nsid: tangled.PipelineNSID, EventJson: eventJson, } return d.insertEvent(event, n)}
// the envelope Created stays zero so insertEvent stamps the local clock, the record's CreatedAt is separatefunc statusEvent(pipelineAtUri, workflow, status string, workflowError *string, exitCode *int64) (eventstream.Event, error) { s := tangled.PipelineStatus{ CreatedAt: time.Now().Format(time.RFC3339), Error: workflowError, ExitCode: exitCode, Pipeline: pipelineAtUri, Workflow: workflow, Status: status, }
eventJson, err := json.Marshal(s) if err != nil { return eventstream.Event{}, err }
return eventstream.Event{ Rkey: tid.TID(), Nsid: tangled.PipelineStatusNSID, EventJson: eventJson, }, nil}
func (d *DB) createStatusEvent( workflowId models.WorkflowId, statusKind models.StatusKind, workflowError *string, exitCode *int64, n *notifier.Notifier,) error { event, err := statusEvent(string(workflowId.PipelineId.AtUri()), workflowId.Name, string(statusKind), workflowError, exitCode) if err != nil { return err } return d.insertEvent(event, n)}
// stamps the mill's own clock so it orders against the cursor like a local writefunc (d *DB) InsertEventStatus( pipelineAtUri string, workflow string, status string, workflowError *string, exitCode *int64, n *notifier.Notifier,) error { event, err := statusEvent(pipelineAtUri, workflow, status, workflowError, exitCode) if err != nil { return err } return d.insertEvent(event, n)}
// deleting the lease in the same transaction prevents the terminal event// from replayingfunc (d *DB) CompleteMillLease( leaseID string, pipelineAtUri string, workflow string, status string, workflowError *string, exitCode *int64, n *notifier.Notifier,) error { return d.ApplyEventBatch(n, func(tx *EventBatchTx) error { if err := tx.InsertStatusEvent(pipelineAtUri, workflow, status, workflowError, exitCode); err != nil { return err } return tx.DeleteLease(leaseID) })}
func (d *DB) GetStatus(workflowId models.WorkflowId) (*tangled.PipelineStatus, error) { pipelineAtUri := workflowId.PipelineId.AtUri()
var eventJson string err := d.QueryRow( ` select event from events where nsid = ? and json_extract(event, '$.pipeline') = ? and json_extract(event, '$.workflow') = ? order by created desc limit 1 `, tangled.PipelineStatusNSID, string(pipelineAtUri), workflowId.Name, ).Scan(&eventJson)
if err != nil { return nil, err }
var status tangled.PipelineStatus if err := json.Unmarshal([]byte(eventJson), &status); err != nil { return nil, err }
return &status, nil}
type statusQueryer interface { QueryRowContext(ctx context.Context, query string, args ...any) *sql.Row}
func workflowStartupDelay(ctx context.Context, q statusQueryer, pipelineAtURI, workflow string) (time.Duration, bool, error) { var pending, running sql.NullInt64 err := q.QueryRowContext(ctx, ` select min(case when json_extract(event, '$.status') = 'pending' then created end), min(case when json_extract(event, '$.status') = 'running' then created end) from events where nsid = 'sh.tangled.pipeline.status' and json_extract(event, '$.pipeline') = ? and json_extract(event, '$.workflow') = ? `, pipelineAtURI, workflow).Scan(&pending, &running) if err != nil { return 0, false, err } if !pending.Valid || !running.Valid { return 0, false, nil } delay := time.Duration(running.Int64 - pending.Int64) if delay < 0 { delay = 0 } return delay, true, nil}
func (d *DB) WorkflowStartupDelay(ctx context.Context, workflowID models.WorkflowId) (time.Duration, bool, error) { return workflowStartupDelay(ctx, d, string(workflowID.PipelineId.AtUri()), workflowID.Name)}
func (tx *EventBatchTx) WorkflowStartupDelay(ctx context.Context, pipelineAtURI, workflow string) (time.Duration, bool, error) { return workflowStartupDelay(ctx, tx.tx, pipelineAtURI, workflow)}
func (tx *EventBatchTx) HasWorkflowStatus(ctx context.Context, pipelineAtURI, workflow, status string) (bool, error) { var present bool err := tx.tx.QueryRowContext(ctx, ` select exists( select 1 from events where nsid = 'sh.tangled.pipeline.status' and json_extract(event, '$.pipeline') = ? and json_extract(event, '$.workflow') = ? and json_extract(event, '$.status') = ? ) `, pipelineAtURI, workflow, status).Scan(&present) return present, err}
func (d *DB) StatusPending(workflowId models.WorkflowId, n *notifier.Notifier) error { return d.createStatusEvent(workflowId, models.StatusKindPending, nil, nil, n)}
func (d *DB) StatusRunning(workflowId models.WorkflowId, n *notifier.Notifier) error { return d.createStatusEvent(workflowId, models.StatusKindRunning, nil, nil, n)}
func (d *DB) StatusFailed(workflowId models.WorkflowId, workflowError string, exitCode int64, n *notifier.Notifier) error { return d.createStatusEvent(workflowId, models.StatusKindFailed, &workflowError, &exitCode, n)}
func (d *DB) StatusCancelled(workflowId models.WorkflowId, workflowError string, exitCode int64, n *notifier.Notifier) error { return d.createStatusEvent(workflowId, models.StatusKindCancelled, &workflowError, &exitCode, n)}
func (d *DB) StatusSuccess(workflowId models.WorkflowId, n *notifier.Notifier) error { return d.createStatusEvent(workflowId, models.StatusKindSuccess, nil, nil, n)}
func (d *DB) StatusTimeout(workflowId models.WorkflowId, n *notifier.Notifier) error { return d.createStatusEvent(workflowId, models.StatusKindTimeout, nil, nil, n)}
type PipelineWorkflow struct { Knot string Rkey string Name string}
func (d *DB) ListPipelineWorkflows(repoDid string) ([]PipelineWorkflow, error) { rows, err := d.Query( `select rkey, event from events where nsid = 'sh.tangled.pipeline' and coalesce(json_extract(event, '$.triggerMetadata.repo.repoDid'), json_extract(event, '$.triggerMetadata.repo.did')) = ?`, repoDid, ) if err != nil { return nil, err } defer rows.Close()
var out []PipelineWorkflow for rows.Next() { var rkey, raw string if err := rows.Scan(&rkey, &raw); err != nil { return nil, err } var p tangled.Pipeline if err := json.Unmarshal([]byte(raw), &p); err != nil { continue } knot := "" if p.TriggerMetadata != nil && p.TriggerMetadata.Repo != nil { knot = p.TriggerMetadata.Repo.Knot } for _, wf := range p.Workflows { out = append(out, PipelineWorkflow{Knot: knot, Rkey: rkey, Name: wf.Name}) } } return out, rows.Err()}
// PipelineKey is one exact pipeline at-uri identitytype PipelineKey struct { Knot string Rkey string}
func (d *DB) DeleteEventsByRepo(repoDid string, extra []PipelineKey) error { tx, err := d.Begin() if err != nil { return err } defer tx.Rollback()
// use knots recorded in pipeline events, not the repo's current knots rows, err := tx.Query( `select coalesce(json_extract(event, '$.triggerMetadata.repo.knot'), ''), rkey from events where nsid = 'sh.tangled.pipeline' and coalesce(json_extract(event, '$.triggerMetadata.repo.repoDid'), json_extract(event, '$.triggerMetadata.repo.did')) = ?`, repoDid, ) if err != nil { return err } seen := map[PipelineKey]bool{} var keys []PipelineKey add := func(knot, rkey string) { k := PipelineKey{Knot: knot, Rkey: rkey} if k.Knot != "" && k.Rkey != "" && !seen[k] { seen[k] = true keys = append(keys, k) } } for rows.Next() { var knot, rkey string if err := rows.Scan(&knot, &rkey); err != nil { rows.Close() return err } add(knot, rkey) } if err := rows.Err(); err != nil { rows.Close() return err } rows.Close() for _, k := range extra { add(k.Knot, k.Rkey) }
for _, k := range keys { if _, err := tx.Exec( `delete from events where nsid = 'sh.tangled.pipeline.status' and json_extract(event, '$.pipeline') = ?`, fmt.Sprintf("at://did:web:%s/sh.tangled.pipeline/%s", k.Knot, k.Rkey), ); err != nil { return err } } if _, err := tx.Exec( `delete from events where nsid = 'sh.tangled.pipeline' and coalesce(json_extract(event, '$.triggerMetadata.repo.repoDid'), json_extract(event, '$.triggerMetadata.repo.did')) = ?`, repoDid, ); err != nil { return err } return tx.Commit()}