package 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 separate func 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 write func (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 replaying func (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 identity type 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() }