package 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 state func 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 executors func 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: "", 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 --token-file ") } 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: "", 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 [--to ]") } 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: "", 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 ") } 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: "", 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: "", 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, err := syntax.ParseDID(cmd.Args().Get(0)) if err != nil { return "", fmt.Errorf("invalid did: %w", err) } return did, nil }