From d643cad86bd3d5b0c12cd30bf6d00d4df56929e1 Mon Sep 17 00:00:00 2001 From: Anirudh Oppiliappan Date: Wed, 11 Jun 2025 17:18:28 +0300 Subject: [PATCH] spindle/{db,engine}: rework db to use rkey Signed-off-by: Anirudh Oppiliappan --- spindle/config/config.go | 1 + spindle/db/db.go | 7 ++-- spindle/db/pipelines.go | 90 +++++++++++++++++++++------------------- spindle/engine/engine.go | 4 +- spindle/server.go | 33 ++++++++------- spindle/stream.go | 5 ++- spindle/tid.go | 9 ++++ 7 files changed, 86 insertions(+), 63 deletions(-) create mode 100644 spindle/tid.go diff --git a/spindle/config/config.go b/spindle/config/config.go index 0a3e3879..2532e34c 100644 --- a/spindle/config/config.go +++ b/spindle/config/config.go @@ -11,6 +11,7 @@ type Server struct { DBPath string `env:"DB_PATH, default=spindle.db"` Hostname string `env:"HOSTNAME, required"` JetstreamEndpoint string `env:"JETSTREAM_ENDPOINT, default=wss://jetstream1.us-west.bsky.network/subscribe"` + Dev bool `env:"DEV, default=false"` } type Config struct { diff --git a/spindle/db/db.go b/spindle/db/db.go index 630a0b08..6a970829 100644 --- a/spindle/db/db.go +++ b/spindle/db/db.go @@ -30,8 +30,9 @@ func Make(dbPath string) (*DB, error) { did text primary key ); - create table if not exists pipelines ( - at_uri text not null, + create table if not exists pipeline_status ( + rkey text not null, + pipeline text not null, status text not null, -- only set if status is 'failed' @@ -42,7 +43,7 @@ func Make(dbPath string) (*DB, error) { updated_at timestamp not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), finished_at timestamp, - primary key (at_uri) + primary key (rkey) ); `) if err != nil { diff --git a/spindle/db/pipelines.go b/spindle/db/pipelines.go index df447c23..d5073f1a 100644 --- a/spindle/db/pipelines.go +++ b/spindle/db/pipelines.go @@ -8,21 +8,21 @@ import ( "tangled.sh/tangled.sh/core/knotserver/notifier" ) -type PipelineStatus string +type PipelineRunStatus string var ( - PipelinePending PipelineStatus = "pending" - PipelineRunning PipelineStatus = "running" - PipelineFailed PipelineStatus = "failed" - PipelineTimeout PipelineStatus = "timeout" - PipelineCancelled PipelineStatus = "cancelled" - PipelineSuccess PipelineStatus = "success" + PipelinePending PipelineRunStatus = "pending" + PipelineRunning PipelineRunStatus = "running" + PipelineFailed PipelineRunStatus = "failed" + PipelineTimeout PipelineRunStatus = "timeout" + PipelineCancelled PipelineRunStatus = "cancelled" + PipelineSuccess PipelineRunStatus = "success" ) -type Pipeline struct { - Rkey string `json:"rkey"` - Knot string `json:"knot"` - Status PipelineStatus `json:"status"` +type PipelineStatus struct { + Rkey string `json:"rkey"` + Pipeline string `json:"pipeline"` + Status PipelineRunStatus `json:"status"` // only if Failed Error string `json:"error"` @@ -33,13 +33,14 @@ type Pipeline struct { FinishedAt time.Time `json:"finished_at"` } -func (p Pipeline) AsRecord() *tangled.PipelineStatus { +func (p PipelineStatus) AsRecord() *tangled.PipelineStatus { exitCode64 := int64(p.ExitCode) finishedAt := p.FinishedAt.String() return &tangled.PipelineStatus{ - Pipeline: fmt.Sprintf("at://%s/%s", p.Knot, p.Rkey), - Status: string(p.Status), + LexiconTypeID: tangled.PipelineStatusNSID, + Pipeline: p.Pipeline, + Status: string(p.Status), ExitCode: &exitCode64, Error: &p.Error, @@ -54,11 +55,11 @@ func pipelineAtUri(rkey, knot string) string { return fmt.Sprintf("at://%s/did:web:%s/%s", tangled.PipelineStatusNSID, knot, rkey) } -func (db *DB) CreatePipeline(rkey, knot string, n *notifier.Notifier) error { +func (db *DB) CreatePipeline(rkey, pipeline string, n *notifier.Notifier) error { _, err := db.Exec(` - insert into pipelines (at_uri, status) - values (?, ?) - `, pipelineAtUri(rkey, knot), PipelinePending) + insert into pipeline_status (rkey, status, pipeline) + values (?, ?, ?) + `, rkey, PipelinePending, pipeline) if err != nil { return err @@ -67,12 +68,12 @@ func (db *DB) CreatePipeline(rkey, knot string, n *notifier.Notifier) error { return nil } -func (db *DB) MarkPipelineRunning(rkey, knot string, n *notifier.Notifier) error { +func (db *DB) MarkPipelineRunning(rkey string, n *notifier.Notifier) error { _, err := db.Exec(` - update pipelines + update pipeline_status set status = ?, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') - where at_uri = ? - `, PipelineRunning, pipelineAtUri(rkey, knot)) + where rkey = ? + `, PipelineRunning, rkey) if err != nil { return err @@ -81,16 +82,16 @@ func (db *DB) MarkPipelineRunning(rkey, knot string, n *notifier.Notifier) error return nil } -func (db *DB) MarkPipelineFailed(rkey, knot string, exitCode int, errorMsg string, n *notifier.Notifier) error { +func (db *DB) MarkPipelineFailed(rkey string, exitCode int, errorMsg string, n *notifier.Notifier) error { _, err := db.Exec(` - update pipelines + update pipeline_status set status = ?, exit_code = ?, error = ?, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now'), finished_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') - where at_uri = ? - `, PipelineFailed, exitCode, errorMsg, pipelineAtUri(rkey, knot)) + where rkey = ? + `, PipelineFailed, exitCode, errorMsg, rkey) if err != nil { return err } @@ -98,12 +99,12 @@ func (db *DB) MarkPipelineFailed(rkey, knot string, exitCode int, errorMsg strin return nil } -func (db *DB) MarkPipelineTimeout(rkey, knot string, n *notifier.Notifier) error { +func (db *DB) MarkPipelineTimeout(rkey string, n *notifier.Notifier) error { _, err := db.Exec(` - update pipelines + update pipeline_status set status = ?, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') - where at_uri = ? - `, PipelineTimeout, pipelineAtUri(rkey, knot)) + where rkey = ? + `, PipelineTimeout, rkey) if err != nil { return err } @@ -111,13 +112,13 @@ func (db *DB) MarkPipelineTimeout(rkey, knot string, n *notifier.Notifier) error return nil } -func (db *DB) MarkPipelineSuccess(rkey, knot string, n *notifier.Notifier) error { +func (db *DB) MarkPipelineSuccess(rkey string, n *notifier.Notifier) error { _, err := db.Exec(` - update pipelines + update pipeline_status set status = ?, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now'), finished_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') - where at_uri = ? - `, PipelineSuccess, pipelineAtUri(rkey, knot)) + where rkey = ? + `, PipelineSuccess, rkey) if err != nil { return err @@ -126,17 +127,17 @@ func (db *DB) MarkPipelineSuccess(rkey, knot string, n *notifier.Notifier) error return nil } -func (db *DB) GetPipeline(rkey, knot string) (Pipeline, error) { - var p Pipeline +func (db *DB) GetPipelineStatus(rkey string) (PipelineStatus, error) { + var p PipelineStatus err := db.QueryRow(` select rkey, status, error, exit_code, started_at, updated_at, finished_at from pipelines - where at_uri = ? - `, pipelineAtUri(rkey, knot)).Scan(&p.Rkey, &p.Status, &p.Error, &p.ExitCode, &p.StartedAt, &p.UpdatedAt, &p.FinishedAt) + where rkey = ? + `, rkey).Scan(&p.Rkey, &p.Status, &p.Error, &p.ExitCode, &p.StartedAt, &p.UpdatedAt, &p.FinishedAt) return p, err } -func (db *DB) GetPipelines(cursor string) ([]Pipeline, error) { +func (db *DB) GetPipelineStatusAsRecords(cursor string) ([]PipelineStatus, error) { whereClause := "" args := []any{} if cursor != "" { @@ -146,7 +147,7 @@ func (db *DB) GetPipelines(cursor string) ([]Pipeline, error) { query := fmt.Sprintf(` select rkey, status, error, exit_code, started_at, updated_at, finished_at - from pipelines + from pipeline_status %s order by rkey asc limit 100 @@ -158,9 +159,9 @@ func (db *DB) GetPipelines(cursor string) ([]Pipeline, error) { } defer rows.Close() - var pipelines []Pipeline + var pipelines []PipelineStatus for rows.Next() { - var p Pipeline + var p PipelineStatus rows.Scan(&p.Rkey, &p.Status, &p.Error, &p.ExitCode, &p.StartedAt, &p.UpdatedAt, &p.FinishedAt) pipelines = append(pipelines, p) } @@ -169,5 +170,10 @@ func (db *DB) GetPipelines(cursor string) ([]Pipeline, error) { return nil, err } + records := []*tangled.PipelineStatus{} + for _, p := range pipelines { + records = append(records, p.AsRecord()) + } + return pipelines, nil } diff --git a/spindle/engine/engine.go b/spindle/engine/engine.go index 2ca9b09c..29e2316b 100644 --- a/spindle/engine/engine.go +++ b/spindle/engine/engine.go @@ -47,7 +47,7 @@ func New(ctx context.Context, db *db.DB, n *notifier.Notifier) (*Engine, error) // SetupPipeline sets up a new network for the pipeline, and possibly volumes etc. // in the future. In here also goes other setup steps. -func (e *Engine) SetupPipeline(ctx context.Context, pipeline *tangled.Pipeline, id string) error { +func (e *Engine) SetupPipeline(ctx context.Context, pipeline *tangled.Pipeline, atUri, id string) error { e.l.Info("setting up pipeline", "pipeline", id) _, err := e.docker.VolumeCreate(ctx, volume.CreateOptions{ @@ -73,7 +73,7 @@ func (e *Engine) SetupPipeline(ctx context.Context, pipeline *tangled.Pipeline, return err } - err = e.db.CreatePipeline(id, e.n) + err = e.db.CreatePipeline(id, atUri, e.n) return err } diff --git a/spindle/server.go b/spindle/server.go index 168ec7a5..f2beb4d7 100644 --- a/spindle/server.go +++ b/spindle/server.go @@ -70,13 +70,14 @@ func Run(ctx context.Context) error { go func() { logger.Info("starting event consumer") knotEventSource := knotclient.NewEventSource("localhost:5555") - ccfg := knotclient.ConsumerConfig{ - Logger: logger, - ProcessFunc: spindle.exec, - } + + ccfg := knotclient.NewConsumerConfig() + ccfg.Logger = logger + ccfg.Dev = cfg.Server.Dev + ccfg.ProcessFunc = spindle.exec ccfg.AddEventSource(knotEventSource) - ec := knotclient.NewEventConsumer(ccfg) + ec := knotclient.NewEventConsumer(*ccfg) ec.Start(ctx) }() @@ -95,19 +96,23 @@ func (s *Spindle) Router() http.Handler { } func (s *Spindle) exec(ctx context.Context, src knotclient.EventSource, msg knotclient.Message) error { - pipeline := tangled.Pipeline{} - err := json.Unmarshal(msg.EventJson, &pipeline) - if err != nil { - fmt.Println("error unmarshalling", err) - return err - } - if msg.Nsid == tangled.PipelineNSID { - err = s.eng.SetupPipeline(ctx, &pipeline, msg.Rkey) + pipeline := tangled.Pipeline{} + err := json.Unmarshal(msg.EventJson, &pipeline) + if err != nil { + fmt.Println("error unmarshalling", err) + return err + } + + // this is a "fake" at uri for now + pipelineAtUri := fmt.Sprintf("at://%s/did:web:%s/%s", tangled.PipelineNSID, pipeline.TriggerMetadata.Repo.Knot, msg.Rkey) + + rkey := TID() + err = s.eng.SetupPipeline(ctx, &pipeline, pipelineAtUri, rkey) if err != nil { return err } - err = s.eng.StartWorkflows(ctx, &pipeline, msg.Rkey) + err = s.eng.StartWorkflows(ctx, &pipeline, rkey) if err != nil { return err } diff --git a/spindle/stream.go b/spindle/stream.go index a1e33bab..9167d47f 100644 --- a/spindle/stream.go +++ b/spindle/stream.go @@ -4,8 +4,9 @@ import ( "net/http" "time" + "context" + "github.com/gorilla/websocket" - "golang.org/x/net/context" ) var upgrader = websocket.Upgrader{ @@ -74,7 +75,7 @@ func (s *Spindle) Events(w http.ResponseWriter, r *http.Request) { } func (s *Spindle) streamPipelines(conn *websocket.Conn, cursor *string) error { - ops, err := s.db.GetPipelines(*cursor) + ops, err := s.db.GetPipelineStatusAsRecords(*cursor) if err != nil { s.l.Debug("err", "err", err) return err diff --git a/spindle/tid.go b/spindle/tid.go new file mode 100644 index 00000000..97802e9b --- /dev/null +++ b/spindle/tid.go @@ -0,0 +1,9 @@ +package spindle + +import "github.com/bluesky-social/indigo/atproto/syntax" + +var TIDClock = syntax.NewTIDClock(0) + +func TID() string { + return TIDClock.Next().String() +} -- 2.51.2