This repository has no description
0

Configure Feed

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

core / appview / state / spindlestream.go
3.2 kB 118 lines
1package state 2 3import ( 4 "context" 5 "encoding/json" 6 "fmt" 7 "log/slog" 8 "strings" 9 "time" 10 11 "github.com/bluesky-social/indigo/atproto/syntax" 12 "tangled.org/core/api/tangled" 13 "tangled.org/core/appview/cache" 14 "tangled.org/core/appview/config" 15 "tangled.org/core/appview/db" 16 "tangled.org/core/appview/models" 17 "tangled.org/core/appview/pipelines" 18 ec "tangled.org/core/eventconsumer" 19 "tangled.org/core/eventconsumer/cursor" 20 "tangled.org/core/log" 21 "tangled.org/core/orm" 22 "tangled.org/core/rbac" 23 spindle "tangled.org/core/spindle/models" 24) 25 26func Spindlestream(ctx context.Context, c *config.Config, d *db.DB, enforcer *rbac.Enforcer, pn *pipelines.StatusNotifier) (*ec.Consumer, error) { 27 logger := log.FromContext(ctx) 28 logger = log.SubLogger(logger, "spindlestream") 29 30 spindles, err := db.GetSpindles( 31 ctx, 32 d, 33 orm.FilterIsNot("verified", "null"), 34 ) 35 if err != nil { 36 return nil, err 37 } 38 39 srcs := make(map[ec.Source]struct{}) 40 for _, s := range spindles { 41 src := ec.NewSpindleSource(s.Instance) 42 srcs[src] = struct{}{} 43 } 44 45 cache := cache.New(c.Redis.Addr) 46 cursorStore := cursor.NewRedisCursorStore(cache) 47 48 cfg := ec.ConsumerConfig{ 49 Sources: srcs, 50 ProcessFunc: spindleIngester(ctx, logger, d, pn), 51 RetryInterval: c.Spindlestream.RetryInterval, 52 MaxRetryInterval: c.Spindlestream.MaxRetryInterval, 53 ConnectionTimeout: c.Spindlestream.ConnectionTimeout, 54 WorkerCount: c.Spindlestream.WorkerCount, 55 QueueSize: c.Spindlestream.QueueSize, 56 Logger: logger, 57 Dev: c.Core.Dev, 58 CursorStore: &cursorStore, 59 } 60 61 return ec.NewConsumer(cfg), nil 62} 63 64func spindleIngester(ctx context.Context, logger *slog.Logger, d *db.DB, pn *pipelines.StatusNotifier) ec.ProcessFunc { 65 return func(ctx context.Context, source ec.Source, msg ec.Message) error { 66 switch msg.Nsid { 67 case tangled.PipelineStatusNSID: 68 return ingestPipelineStatus(ctx, logger, d, pn, source, msg) 69 } 70 71 return nil 72 } 73} 74 75func ingestPipelineStatus(ctx context.Context, logger *slog.Logger, d *db.DB, pn *pipelines.StatusNotifier, source ec.Source, msg ec.Message) error { 76 var record tangled.PipelineStatus 77 err := json.Unmarshal(msg.EventJson, &record) 78 if err != nil { 79 return err 80 } 81 82 pipelineUri, err := syntax.ParseATURI(record.Pipeline) 83 if err != nil { 84 return err 85 } 86 87 exitCode := 0 88 if record.ExitCode != nil { 89 exitCode = int(*record.ExitCode) 90 } 91 92 // pick the record creation time if possible, or use time.Now 93 created := time.Now() 94 if t, err := time.Parse(time.RFC3339, record.CreatedAt); err == nil && created.After(t) { 95 created = t 96 } 97 98 status := models.PipelineStatus{ 99 Spindle: source.Key(), 100 Rkey: msg.Rkey, 101 PipelineKnot: strings.TrimPrefix(pipelineUri.Authority().String(), "did:web:"), 102 PipelineRkey: pipelineUri.RecordKey().String(), 103 Created: created, 104 Workflow: record.Workflow, 105 Status: spindle.StatusKind(record.Status), 106 Error: record.Error, 107 ExitCode: exitCode, 108 } 109 110 err = db.AddPipelineStatus(ctx, d, status) 111 if err != nil { 112 return fmt.Errorf("failed to add pipeline status: %w", err) 113 } 114 115 pn.Publish(pipelineUri) 116 117 return nil 118}