Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648package main
import ( "context" "crypto/rand" "encoding/hex" "fmt" "log/slog" "math" "os" "os/signal" "strings" "syscall" "text/tabwriter" "time"
"github.com/bluesky-social/indigo/atproto/syntax" "github.com/urfave/cli/v3"
tlog "tangled.org/core/log" "tangled.org/core/spindle" "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" "tangled.org/core/spindle/mill" "tangled.org/core/spindle/quota")
func main() { cmd := &cli.Command{ Name: "spindle", Usage: "spindle continuous integration runner", Commands: []*cli.Command{ Command(), millCommand(), executorCommand(), quotaCommand(), banCommand(), unbanCommand(), bansCommand(), }, DefaultCommand: "run", }
format := os.Getenv("SPINDLE_LOG_FORMAT") logger, err := tlog.NewWithFormat("spindle", format) if err != nil { fmt.Fprintf(os.Stderr, "invalid log format: %v\n", err) os.Exit(1) } slog.SetDefault(logger)
ctx := tlog.IntoContext(context.Background(), logger) ctx, stopSignals := signal.NotifyContext(ctx, os.Interrupt, syscall.SIGTERM) defer stopSignals()
if err := cmd.Run(ctx, os.Args); err != nil { logger.Error(err.Error()) os.Exit(-1) }}
func Command() *cli.Command { return &cli.Command{ Name: "run", Usage: "run the spindle server", Action: func(ctx context.Context, cmd *cli.Command) error { return spindle.Run(ctx) }, }}
var dbFlag = &cli.StringFlag{ Name: "db", Usage: "path to the spindle sqlite db", Value: "spindle.db", Sources: cli.EnvVars("SPINDLE_SERVER_DB_PATH"),}
func openDB(ctx context.Context, cmd *cli.Command) (*db.DB, error) { return db.Make(ctx, cmd.String("db"))}
// lets executor operators reset stream statefunc executorCommand() *cli.Command { return &cli.Command{ Name: "executor", Usage: "executor host administration", Commands: []*cli.Command{ { Name: "reset-stream", Usage: "wipe the outbox and start a fresh stream epoch, for when the mill has lost this executor's acked position; restart the executor service afterwards", Flags: []cli.Flag{dbFlag}, Action: func(ctx context.Context, cmd *cli.Command) error { d, err := openDB(ctx, cmd) if err != nil { return err } var b [8]byte if _, err := rand.Read(b[:]); err != nil { return err } epoch := hex.EncodeToString(b[:]) if err := d.SetOutboxEpoch(epoch); err != nil { return fmt.Errorf("resetting outbox: %w", err) } // clear artifacts waiting on leases from the old // stream, replaying them under a new epoch would // trip the mill's lease-epoch check if err := d.ClearPendingArtifacts(); err != nil { return fmt.Errorf("clearing pending artifacts: %w", err) } fmt.Printf("outbox wiped; new stream epoch %s. restart the executor to reconnect cleanly\n", epoch) return nil }, }, }, }}
// lets mill operators manage executorsfunc millCommand() *cli.Command { return &cli.Command{ Name: "mill", Usage: "mill host administration", Commands: []*cli.Command{ { Name: "executor", Usage: "manage executors allowed to join this mill", Commands: []*cli.Command{ { Name: "token", Usage: "manage executor credentials", Commands: []*cli.Command{ { Name: "generate", Usage: "generate and store an executor token", Flags: []cli.Flag{ dbFlag, &cli.DurationFlag{ Name: "ttl", Usage: "token lifetime (e.g. 720h); omit for no expiry", }, }, Action: func(ctx context.Context, cmd *cli.Command) error { d, err := openDB(ctx, cmd) if err != nil { return err } token, err := mill.GenerateToken() if err != nil { return err } var expiresAt *time.Time if ttl := cmd.Duration("ttl"); ttl > 0 { exp := time.Now().UTC().Add(ttl) expiresAt = &exp } if err := d.CreateExecutorToken(mill.HashToken(token), expiresAt); err != nil { return fmt.Errorf("storing executor token: %w", err) } fmt.Println(token) return nil }, }, }, }, { Name: "add", Usage: "register an executor with an existing token", ArgsUsage: "<name>", Flags: []cli.Flag{ dbFlag, &cli.StringFlag{ Name: "token-file", Usage: "file containing a token created by executor token generate", }, &cli.StringSliceFlag{ Name: "label", Usage: "authorized labels for this executor", }, }, Action: func(ctx context.Context, cmd *cli.Command) error { name := cmd.Args().First() if name == "" { return fmt.Errorf("usage: spindle mill executor add <name> --token-file <path>") } tokenFile := cmd.String("token-file") if tokenFile == "" { return fmt.Errorf("--token-file is required") } tokenBytes, err := os.ReadFile(tokenFile) if err != nil { return fmt.Errorf("reading token file: %w", err) } token := strings.TrimSpace(string(tokenBytes)) if token == "" { return fmt.Errorf("token file is empty") } d, err := openDB(ctx, cmd) if err != nil { return err } labels := cmd.StringSlice("label") if err := d.RegisterExecutor(name, mill.HashToken(token), labels); err != nil { return fmt.Errorf("registering executor %q: %w", name, err) } return nil }, }, { Name: "list", Usage: "list registered executors", Flags: []cli.Flag{dbFlag}, Action: func(ctx context.Context, cmd *cli.Command) error { d, err := openDB(ctx, cmd) if err != nil { return err } registrations, err := d.ListExecutorRegistrations() if err != nil { return err } w := tabwriter.NewWriter(os.Stdout, 0, 0, 3, ' ', 0) fmt.Fprintln(w, "NAME\tCREATED\tEXPIRES\tLABELS") for _, registration := range registrations { expires := "never" if registration.ExpiresAt != nil { expires = registration.ExpiresAt.Format(time.RFC3339) if time.Now().After(*registration.ExpiresAt) { expires += " (expired)" } } labels := strings.Join(registration.Labels, ",") if labels == "" { labels = "-" } fmt.Fprintf(w, "%s\t%s\t%s\t%s\n", registration.Name, registration.CreatedAt, expires, labels) } return w.Flush() }, }, { Name: "reset-cursor", Usage: "reset the mill's acked stream position for an executor after mill state loss; takes effect on the executor's next reconnect", ArgsUsage: "<name>", Flags: []cli.Flag{ dbFlag, &cli.UintFlag{ Name: "to", Usage: "skip forward to this acked seqno instead of forgetting everything; " + "events at or below it are treated as applied (accepts their loss)", }, }, Action: func(ctx context.Context, cmd *cli.Command) error { name := cmd.Args().First() if name == "" { return fmt.Errorf("usage: spindle mill executor reset-cursor <name> [--to <seqno>]") } d, err := openDB(ctx, cmd) if err != nil { return err } if cmd.IsSet("to") { n, err := d.SetExecutorCursors(name, uint64(cmd.Uint("to"))) if err != nil { return err } if n == 0 { return fmt.Errorf("no cursor rows for executor %q", name) } fmt.Printf("skipped %s forward to seqno %d (%d epoch rows); earlier events are lost\n", name, cmd.Uint("to"), n) return nil } n, err := d.DeleteExecutorCursors(name) if err != nil { return err } fmt.Printf("forgot %d cursor row(s) for %s; the mill now expects its stream from seqno 1\n", n, name) return nil }, }, { Name: "revoke", Usage: "revoke an executor's token", ArgsUsage: "<name>", Flags: []cli.Flag{dbFlag}, Action: func(ctx context.Context, cmd *cli.Command) error { name := cmd.Args().First() if name == "" { return fmt.Errorf("usage: spindle mill executor revoke <name>") } d, err := openDB(ctx, cmd) if err != nil { return err } ok, err := d.RevokeExecutorToken(name) if err != nil { return err } if !ok { return fmt.Errorf("no such executor identity %q", name) } return nil }, }, }, }, }, }}
func openQuotaStore(ctx context.Context, cmd *cli.Command) (*db.QuotaStore, error) { d, err := openDB(ctx, cmd) if err != nil { return nil, err } return db.NewQuotaStore(d, quota.Defaults{}), nil}
type quotaKey struct { did string resource string}
func quotaCommand() *cli.Command { return &cli.Command{ Name: "quota", Usage: "manage resource quotas for repos and users", Commands: []*cli.Command{ { Name: "usage", Usage: "show usage against limits for every repo and user", Flags: []cli.Flag{dbFlag}, Action: func(ctx context.Context, cmd *cli.Command) error { if cmd.Args().Len() != 0 { return fmt.Errorf("unexpected arguments") } qs, err := openQuotaStore(ctx, cmd) if err != nil { return err } defer qs.Close()
defaults, err := config.LoadQuotaDefaults(ctx) if err != nil { return fmt.Errorf("invalid quota configuration: %w", err) }
limits, err := qs.ListLimits(ctx) if err != nil { return fmt.Errorf("listing limits: %w", err) } overrides := make(map[quotaKey]int64, len(limits)) for _, l := range limits { overrides[quotaKey{did: l.DID, resource: l.Resource}] = l.Limit }
usages, err := qs.ListUsage(ctx) if err != nil { return fmt.Errorf("listing usage: %w", err) }
w := tabwriter.NewWriter(os.Stdout, 0, 0, 3, ' ', 0) fmt.Fprintln(w, "SCOPE\tDID\tRESOURCE\tUSED\tLIMIT")
for _, u := range usages { var limitVal int64 = -1 if ovLimit, ok := overrides[quotaKey{did: u.DID, resource: u.Resource}]; ok { limitVal = ovLimit } else if scopeDefs, ok := defaults[u.Scope]; ok { if defLimit, ok := scopeDefs[u.Resource]; ok { limitVal = normalizeDefaultLimit(defLimit) } } fmt.Fprintf(w, "%s\t%s\t%s\t%s\t%s\n", string(u.Scope), u.DID, string(u.Resource), formatLimit(u.Resource, u.Used), formatLimit(u.Resource, limitVal)) } return w.Flush() }, }, { Name: "list", Usage: "list all custom limits", Flags: []cli.Flag{dbFlag}, Action: func(ctx context.Context, cmd *cli.Command) error { if cmd.Args().Len() != 0 { return fmt.Errorf("unexpected arguments") } qs, err := openQuotaStore(ctx, cmd) if err != nil { return err } defer qs.Close() limits, err := qs.ListLimits(ctx) if err != nil { return fmt.Errorf("listing limits: %w", err) } w := tabwriter.NewWriter(os.Stdout, 0, 0, 3, ' ', 0) fmt.Fprintln(w, "DID\tRESOURCE\tLIMIT") for _, l := range limits { fmt.Fprintf(w, "%s\t%s\t%s\n", l.DID, string(l.Resource), formatLimit(l.Resource, l.Limit)) } return w.Flush() }, }, { Name: "set", Usage: "give a repo or user a custom limit", Flags: []cli.Flag{ dbFlag, &cli.StringFlag{ Name: "did", Usage: "did of the user or repo", Required: true, }, &cli.StringFlag{ Name: "resource", Usage: "resource to limit (e.g. workflows, cache_storage_bytes)", Required: true, }, &cli.IntFlag{ Name: "limit", Usage: "new limit (MiB for cache_storage_bytes, memory_mib, and disk_mib)", }, &cli.BoolFlag{ Name: "unlimited", Usage: "remove the limit entirely", }, }, Action: func(ctx context.Context, cmd *cli.Command) error { if cmd.Args().Len() != 0 { return fmt.Errorf("unexpected arguments") }
did, err := syntax.ParseDID(cmd.String("did")) if err != nil { return fmt.Errorf("invalid did: %w", err) }
resource := cmd.String("resource")
if err := quota.ValidateOverride(did.String(), resource); err != nil { return err }
hasLimit := cmd.IsSet("limit") hasUnlimited := cmd.Bool("unlimited")
if !hasLimit && !hasUnlimited { return fmt.Errorf("must specify either --limit or --unlimited") } if hasLimit && hasUnlimited { return fmt.Errorf("cannot specify both --limit and --unlimited") }
var limitVal int64 = -1 if hasLimit { limitInt := cmd.Int("limit") if limitInt < 0 { return fmt.Errorf("limit cannot be negative: %d", limitInt) }
if resource == quota.ResourceCacheStorageBytes { if limitInt > math.MaxInt64/1048576 { return fmt.Errorf("limit %d MiB would overflow bytes", limitInt) } limitVal = int64(limitInt) * 1048576 } else { limitVal = int64(limitInt) } }
qs, err := openQuotaStore(ctx, cmd) if err != nil { return err } defer qs.Close() if err := qs.SetLimit(ctx, did.String(), resource, limitVal); err != nil { return fmt.Errorf("setting limit: %w", err) } fmt.Printf("quota override set for %s %s\n", did.String(), resource) return nil }, }, { Name: "unset", Usage: "remove a custom limit and go back to the defaults", Flags: []cli.Flag{ dbFlag, &cli.StringFlag{ Name: "did", Usage: "did of the user or repo", Required: true, }, &cli.StringFlag{ Name: "resource", Usage: "resource to limit (e.g. workflows, cache_storage_bytes)", Required: true, }, }, Action: func(ctx context.Context, cmd *cli.Command) error { if cmd.Args().Len() != 0 { return fmt.Errorf("unexpected arguments") }
did, err := syntax.ParseDID(cmd.String("did")) if err != nil { return fmt.Errorf("invalid did: %w", err) }
resource := cmd.String("resource")
if err := quota.ValidateOverride(did.String(), resource); err != nil { return err }
qs, err := openQuotaStore(ctx, cmd) if err != nil { return err } defer qs.Close() if err := qs.UnsetLimit(ctx, did.String(), resource); err != nil { return fmt.Errorf("unsetting limit: %w", err) } fmt.Printf("quota override removed for %s %s\n", did.String(), resource) return nil }, }, }, }}
func normalizeDefaultLimit(limit int64) int64 { if limit == 0 { return -1 } return limit}
func formatLimit(resource string, limit int64) string { if limit < 0 { return "unlimited" } switch resource { case quota.ResourceCacheStorageBytes: if limit%(1024*1024) == 0 { return fmt.Sprintf("%d MiB", limit/(1024*1024)) } return fmt.Sprintf("%d B", limit) case quota.ResourceMemoryMiB: return fmt.Sprintf("%d MiB", limit) case quota.ResourceDiskMiB: return fmt.Sprintf("%d MiB", limit) case quota.ResourceWorkflows: return fmt.Sprintf("%d workflows", limit) case quota.ResourceVCPUs: return fmt.Sprintf("%d vCPUs", limit) case quota.ResourceWebhooks: return fmt.Sprintf("%d webhooks", limit) default: return fmt.Sprintf("%d", limit) }}
func banCommand() *cli.Command { return &cli.Command{ Name: "ban", ArgsUsage: "<did>", Usage: "ban a repo or user and wipe their existing data on next startup", Flags: []cli.Flag{dbFlag}, Action: func(ctx context.Context, cmd *cli.Command) error { did, err := banArg(cmd) if err != nil { return err } d, err := openDB(ctx, cmd) if err != nil { return err } defer d.Close() return d.PutBan(db.BanEntry{SubjectDid: did}) }, }}
func unbanCommand() *cli.Command { return &cli.Command{ Name: "unban", ArgsUsage: "<did>", Usage: "unban a repo or user", Flags: []cli.Flag{dbFlag}, Action: func(ctx context.Context, cmd *cli.Command) error { did, err := banArg(cmd) if err != nil { return err } d, err := openDB(ctx, cmd) if err != nil { return err } defer d.Close() removed, err := d.DeleteBan(did) if err != nil { return err } if !removed { return fmt.Errorf("%s is not banned", did) } return nil }, }}
func bansCommand() *cli.Command { return &cli.Command{ Name: "bans", Usage: "list all bans", Flags: []cli.Flag{dbFlag}, Action: func(ctx context.Context, cmd *cli.Command) error { d, err := openDB(ctx, cmd) if err != nil { return err } defer d.Close() bans, err := d.BanList() if err != nil { return err } w := tabwriter.NewWriter(os.Stdout, 0, 4, 2, ' ', 0) fmt.Fprintln(w, "DID\tSINCE") for _, b := range bans { fmt.Fprintf(w, "%s\t%s\n", b.SubjectDid, b.CreatedAt) } return w.Flush() }, }}
func banArg(cmd *cli.Command) (syntax.DID, error) { if cmd.Args().Len() != 1 { return "", fmt.Errorf("expected <did>") } did, err := syntax.ParseDID(cmd.Args().Get(0)) if err != nil { return "", fmt.Errorf("invalid did: %w", err) } return did, nil}