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