From 1edcd7707b74263a60a244da3cbb982d5016a0ae Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Tue, 30 Dec 2025 12:35:02 +0900 Subject: [PATCH] wip: spindle: engines -> adapters Signed-off-by: Seongmin Lee --- spindle/adapters/nixery/adapter.go | 459 ++++++++++++++++++++++++++++ spindle/adapters/nixery/readme.md | 3 + spindle/adapters/nixery/workflow.go | 42 +++ spindle/models/adapter.go | 93 ++++++ spindle/models/pipeline2.go | 124 ++++++++ spindle/pipeline.go | 40 +++ spindle/repomanager/repomanager.go | 169 ++++++++++ spindle/server.go | 1 + 8 files changed, 931 insertions(+) create mode 100644 spindle/adapters/nixery/adapter.go create mode 100644 spindle/adapters/nixery/readme.md create mode 100644 spindle/adapters/nixery/workflow.go create mode 100644 spindle/models/adapter.go create mode 100644 spindle/models/pipeline2.go create mode 100644 spindle/pipeline.go create mode 100644 spindle/repomanager/repomanager.go diff --git a/spindle/adapters/nixery/adapter.go b/spindle/adapters/nixery/adapter.go new file mode 100644 index 00000000..4af3b2ef --- /dev/null +++ b/spindle/adapters/nixery/adapter.go @@ -0,0 +1,459 @@ +package nixery + +import ( + "context" + "fmt" + "io" + "log/slog" + "os" + "path" + "path/filepath" + "regexp" + "runtime" + "sync" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/docker/docker/api/types/container" + "github.com/docker/docker/api/types/filters" + "github.com/docker/docker/api/types/image" + "github.com/docker/docker/api/types/mount" + "github.com/docker/docker/api/types/network" + "github.com/docker/docker/client" + "github.com/stretchr/testify/assert/yaml" + "tangled.org/core/api/tangled" + "tangled.org/core/sets" + "tangled.org/core/spindle/config" + "tangled.org/core/spindle/models" + "tangled.org/core/spindle/repomanager" + "tangled.org/core/tid" + "tangled.org/core/workflow" +) + +const AdapterID = "nixery" + +type Adapter struct { + l *slog.Logger + repoManager *repomanager.RepoManager + docker client.APIClient + Timeout time.Duration + spindleDid syntax.DID + cfg config.NixeryPipelines + + mu sync.RWMutex + activeRuns map[syntax.ATURI]models.WorkflowRun + subscribers sets.Set[chan<- models.WorkflowRun] +} + +var _ models.Adapter = (*Adapter)(nil) + +func New(l *slog.Logger, cfg config.Config, repoManager *repomanager.RepoManager) (*Adapter, error) { + dc, err := client.NewClientWithOpts(client.FromEnv, client.WithAPIVersionNegotiation()) + if err != nil { + return nil, fmt.Errorf("creating docker client: %w", err) + } + return &Adapter{ + l: l, + repoManager: repoManager, + docker: dc, + Timeout: time.Minute * 5, // TODO: set timeout from config + spindleDid: cfg.Server.Did(), + cfg: cfg.NixeryPipelines, + + activeRuns: make(map[syntax.ATURI]models.WorkflowRun), + subscribers: sets.New[chan<- models.WorkflowRun](), + }, nil +} + +func (a *Adapter) Init() error { + // no-op + return nil +} + +func (a *Adapter) Shutdown(ctx context.Context) error { + // TODO: cleanup spawned containers just in case + panic("unimplemented") +} + +func (a *Adapter) SetupRepo(ctx context.Context, repo syntax.ATURI) error { + if err := a.repoManager.RegisterRepo(ctx, repo, []string{"/.tangled/workflows"}); err != nil { + return fmt.Errorf("syncing repo: %w", err) + } + return nil +} + +func (a *Adapter) ListWorkflowDefs(ctx context.Context, repo syntax.ATURI, rev string) ([]models.WorkflowDef, error) { + defs, err := a.listWorkflowDefs(ctx, repo, rev) + if err != nil { + return nil, err + } + retDefs := make([]models.WorkflowDef, len(defs)) + for i, def := range defs { + retDefs[i] = def.AsInfo() + } + return retDefs, nil +} + +func (a *Adapter) listWorkflowDefs(ctx context.Context, repo syntax.ATURI, rev string) ([]WorkflowDef, error) { + workflowDir, err := a.repoManager.FileTree(ctx, repo, rev, workflow.WorkflowDir) + if err != nil { + return nil, fmt.Errorf("loading file tree: %w", err) + } + + if len(workflowDir) == 0 { + return nil, nil + } + + // TODO(boltless): repoManager.FileTree() should be smart enough so we don't need to do this: + gr, err := a.repoManager.Open(repo, rev) + if err != nil { + return nil, fmt.Errorf("opening git repo: %w", err) + } + + var defs []WorkflowDef + for _, e := range workflowDir { + if !e.IsFile() { + continue + } + + fpath := filepath.Join(workflow.WorkflowDir, e.Name) + contents, err := gr.RawContent(fpath) + if err != nil { + return nil, fmt.Errorf("reading raw content of '%s': %w", fpath, err) + } + + var wf WorkflowDef + if err := yaml.Unmarshal(contents, &wf); err != nil { + return nil, fmt.Errorf("parsing yaml: %w", err) + } + wf.Name = e.Name + + defs = append(defs, wf) + } + + return defs, nil +} + +func (a *Adapter) EvaluateEvent(ctx context.Context, event models.Event) ([]models.WorkflowRun, error) { + defs, err := a.listWorkflowDefs(ctx, event.SourceRepo, event.SourceSha) + if err != nil { + return nil, fmt.Errorf("fetching workflow definitions: %w", err) + } + + // filter out triggered workflows + var triggered []nixeryWorkflow + for _, def := range defs { + if def.ShouldRunOn(event) { + triggered = append(triggered, nixeryWorkflow{ + event: event, + def: def, + }) + } + } + + // TODO: append more workflows from "on_workflow" event + + // schedule workflows and return immediately + runs := make([]models.WorkflowRun, len(triggered)) + for i, workflow := range triggered { + runs[i] = a.scheduleWorkflow(ctx, workflow) + } + return runs, nil +} + +// NOTE: nixery adapter is volatile. GetActiveWorkflowRun will return error +// when the workflow is terminated. It lets spindle to mark lost workflow.run +// as "Failed". +func (a *Adapter) GetActiveWorkflowRun(ctx context.Context, runId syntax.ATURI) (models.WorkflowRun, error) { + a.mu.RLock() + run, exists := a.activeRuns[runId] + a.mu.RUnlock() + if !exists { + return run, fmt.Errorf("unknown or terminated workflow") + } + return run, nil +} + +func (a *Adapter) ListActiveWorkflowRuns(ctx context.Context) ([]models.WorkflowRun, error) { + a.mu.RLock() + defer a.mu.RUnlock() + + runs := make([]models.WorkflowRun, 0, len(a.activeRuns)) + for _, run := range a.activeRuns { + runs = append(runs, run) + } + return runs, nil +} + +func (a *Adapter) SubscribeWorkflowRun(ctx context.Context) <-chan models.WorkflowRun { + ch := make(chan models.WorkflowRun, 1) + + a.mu.Lock() + a.subscribers.Insert(ch) + a.mu.Unlock() + + // cleanup spindle stops listening + go func() { + <-ctx.Done() + a.mu.Lock() + a.subscribers.Remove(ch) + a.mu.Unlock() + close(ch) + }() + + return ch +} + +func (a *Adapter) emit(run models.WorkflowRun) { + a.mu.Lock() + if run.Status.IsActive() { + a.activeRuns[run.AtUri()] = run + } else { + delete(a.activeRuns, run.AtUri()) + } + + // Snapshot subscribers to broadcast outside the lock + subs := make([]chan<- models.WorkflowRun, 0, a.subscribers.Len()) + for ch := range a.subscribers.All() { + subs = append(subs, ch) + } + a.mu.Unlock() + + for _, ch := range subs { + select { + case ch <- run: + default: + // avoid blocking if channel is full + // spindle will catch the state by regular GetWorkflowRun poll + } + } +} + +func (a *Adapter) StreamWorkflowRunLogs(ctx context.Context, runId syntax.ATURI, handle func(line models.LogLine) error) error { + panic("unimplemented") +} + +func (a *Adapter) CancelWorkflowRun(ctx context.Context, runId syntax.ATURI) error { + // remove network + if err := a.docker.NetworkRemove(ctx, networkName(runId)); err != nil { + return fmt.Errorf("removing network: %w", err) + } + + // stop & remove docker containers with label + containers, err := a.docker.ContainerList(ctx, container.ListOptions{ + Filters: labelFilter(tangled.CiWorkflowRunNSID, runId.String()), + }) + if err != nil { + return fmt.Errorf("finding container with label: %w", err) + } + for _, c := range containers { + if err := a.docker.ContainerStop(ctx, c.ID, container.StopOptions{}); err != nil { + return fmt.Errorf("stopping container: %w", err) + } + + if err := a.docker.ContainerRemove(ctx, c.ID, container.RemoveOptions{ + RemoveVolumes: true, + RemoveLinks: false, + Force: false, + }); err != nil { + return fmt.Errorf("removing container: %w", err) + } + } + return nil +} + +func labelFilter(labelKey, labelVal string) filters.Args { + filterArgs := filters.NewArgs() + filterArgs.Add("label", fmt.Sprintf("%s=%s", labelKey, labelVal)) + return filterArgs +} + +const ( + workspaceDir = "/tangled/workspace" + homeDir = "/tangled/home" +) + +// scheduleWorkflow schedules a workflow run in job queue and return queued run +func (a *Adapter) scheduleWorkflow(ctx context.Context, workflow nixeryWorkflow) models.WorkflowRun { + l := a.l + + run := models.WorkflowRun{ + Did: a.spindleDid, + Rkey: syntax.RecordKey(tid.TID()), + AdapterId: AdapterID, + Name: workflow.def.Name, + Status: models.WorkflowStatusPending, + } + + a.mu.Lock() + a.activeRuns[run.AtUri()] = run + a.mu.Unlock() + + go func() { + defer a.CancelWorkflowRun(ctx, run.AtUri()) + + containerId, err := a.initNixeryContainer(ctx, workflow.def, run.AtUri()) + if err != nil { + l.Error("failed to intialize container", "err", err) + // TODO: put user-facing logs in workflow log + a.emit(run.WithStatus(models.WorkflowStatusFailed)) + return + } + + ctx, cancel := context.WithTimeout(ctx, a.Timeout) + defer cancel() + + for stepIdx, step := range workflow.def.Steps { + if err := a.runStep(ctx, containerId, stepIdx, step); err != nil { + l.Error("failed to run step", "stepIdx", stepIdx, "err", err) + return + } + } + l.Info("all steps completed successfully") + }() + + l.Info("workflow scheduled to background", "workflow.run", run.AtUri()) + + return run +} + +func (a *Adapter) runStep(ctx context.Context, containerId string, stepIdx int, step Step) error { + // TODO: implement this + + // TODO: configure envs + var envs []string + + select { + case <-ctx.Done(): + return ctx.Err() + default: + } + + mkExecResp, err := a.docker.ContainerExecCreate(ctx, containerId, container.ExecOptions{ + Cmd: []string{"bash", "-c", step.Command}, + AttachStdout: true, + AttachStderr: true, + Env: envs, + }) + if err != nil { + return fmt.Errorf("creating exec: %w", err) + } + + panic("unimplemented") +} + +// initNixeryContainer pulls the image from nixery and start the container. +func (a *Adapter) initNixeryContainer(ctx context.Context, def WorkflowDef, runAt syntax.ATURI) (string, error) { + imageName := workflowImageName(def.Dependencies, a.cfg.Nixery) + + _, err := a.docker.NetworkCreate(ctx, networkName(runAt), network.CreateOptions{ + Driver: "bridge", + }) + if err != nil { + return "", fmt.Errorf("creating network: %w", err) + } + + reader, err := a.docker.ImagePull(ctx, imageName, image.PullOptions{}) + if err != nil { + return "", fmt.Errorf("pulling image: %w", err) + } + defer reader.Close() + io.Copy(os.Stdout, reader) + + resp, err := a.docker.ContainerCreate(ctx, &container.Config{ + Image: imageName, + Cmd: []string{"cat"}, + OpenStdin: true, // so cat stays alive :3 + Tty: false, + Hostname: "spindle", + WorkingDir: workspaceDir, + Labels: map[string]string{ + tangled.CiWorkflowRunNSID: runAt.String(), + }, + // TODO(winter): investigate whether environment variables passed here + // get propagated to ContainerExec processes + }, &container.HostConfig{ + Mounts: []mount.Mount{ + { + Type: mount.TypeTmpfs, + Target: "/tmp", + ReadOnly: false, + TmpfsOptions: &mount.TmpfsOptions{ + Mode: 0o1777, // world-writeable sticky bit + Options: [][]string{ + {"exec"}, + }, + }, + }, + }, + ReadonlyRootfs: false, + CapDrop: []string{"ALL"}, + CapAdd: []string{"CAP_DAC_OVERRIDE", "CAP_CHOWN", "CAP_FOWNER", "CAP_SETUID", "CAP_SETGID"}, + SecurityOpt: []string{"no-new-privileges"}, + ExtraHosts: []string{"host.docker.internal:host-gateway"}, + }, nil, nil, "") + if err != nil { + return "", fmt.Errorf("creating container: %w", err) + } + + if err := a.docker.ContainerStart(ctx, resp.ID, container.StartOptions{}); err != nil { + return "", fmt.Errorf("starting container: %w", err) + } + + mkExecResp, err := a.docker.ContainerExecCreate(ctx, resp.ID, container.ExecOptions{ + Cmd: []string{"mkdir", "-p", workspaceDir, homeDir}, + AttachStdout: true, // NOTE(winter): pretty sure this will make it so that when stdout read is done below, mkdir is done. maybe?? + AttachStderr: true, // for good measure, backed up by docker/cli ("If -d is not set, attach to everything by default") + }) + if err != nil { + return "", err + } + + // This actually *starts* the command. Thanks, Docker! + execResp, err := a.docker.ContainerExecAttach(ctx, mkExecResp.ID, container.ExecAttachOptions{}) + if err != nil { + return "", err + } + defer execResp.Close() + + // This is apparently best way to wait for the command to complete. + _, err = io.ReadAll(execResp.Reader) + if err != nil { + return "", err + } + + execInspectResp, err := a.docker.ContainerExecInspect(ctx, mkExecResp.ID) + if err != nil { + return "", err + } + + if execInspectResp.ExitCode != 0 { + return "", fmt.Errorf("mkdir exited with exit code %d", execInspectResp.ExitCode) + } else if execInspectResp.Running { + return "", fmt.Errorf("mkdir is somehow still running??") + } + + return resp.ID, nil +} + +func workflowImageName(deps map[string][]string, nixery string) string { + var dependencies string + for reg, ds := range deps { + if reg == "nixpkgs" { + dependencies = path.Join(ds...) + } + } + // NOTE: shouldn't base dependencies come first? + // like: nixery.tangled.sh/arm64/bash/git/coreutils/nix + dependencies = path.Join(dependencies, "bash", "git", "coreutils", "nix") + if runtime.GOARCH == "arm64" { + dependencies = path.Join("arm64", dependencies) + } + + return path.Join(nixery, dependencies) +} + +var re = regexp.MustCompile(`[^a-zA-Z0-9_.-]`) +func networkName(runId syntax.ATURI) string { + return re.ReplaceAllString(runId.String()[5:], "-") +} diff --git a/spindle/adapters/nixery/readme.md b/spindle/adapters/nixery/readme.md new file mode 100644 index 00000000..cec8a6bc --- /dev/null +++ b/spindle/adapters/nixery/readme.md @@ -0,0 +1,3 @@ +# Nixery spindle adapter implementation + +Nixery adapter uses `/.tangled/workflows/*.yml` files as workflow definitions. diff --git a/spindle/adapters/nixery/workflow.go b/spindle/adapters/nixery/workflow.go new file mode 100644 index 00000000..b62a9c1f --- /dev/null +++ b/spindle/adapters/nixery/workflow.go @@ -0,0 +1,42 @@ +package nixery + +import ( + "tangled.org/core/spindle/models" + "tangled.org/core/workflow" +) + +type nixeryWorkflow struct { + event models.Event // event that triggered the workflow + def WorkflowDef // definition of the workflow +} + +// TODO: extract general fields to workflow.WorkflowDef struct + +// nixery adapter workflow definition spec +type WorkflowDef struct { + Name string `yaml:"-"` // name of the workflow file + When []workflow.Constraint `yaml:"when"` + CloneOpts workflow.CloneOpts `yaml:"clone"` + + Dependencies map[string][]string // nix packages used for the workflow + Steps []Step // workflow steps +} + +type Step struct { + Name string `yaml:"name"` + Command string `yaml:"command"` + Enviornment map[string]string `yaml:"environment"` +} + +func (d *WorkflowDef) AsInfo() models.WorkflowDef { + return models.WorkflowDef{ + AdapterId: AdapterID, + Name: d.Name, + When: d.When, + } +} + +func (d *WorkflowDef) ShouldRunOn(event models.Event) bool { + // panic("unimplemented") + return false +} diff --git a/spindle/models/adapter.go b/spindle/models/adapter.go new file mode 100644 index 00000000..b23f2568 --- /dev/null +++ b/spindle/models/adapter.go @@ -0,0 +1,93 @@ +package models + +import ( + "context" + + "github.com/bluesky-social/indigo/atproto/syntax" +) + +// Adapter is the core of the spindle. It can use its own way to configure and +// run the workflows. The workflow definition can be either yaml files in git +// repositories or even from dedicated web UI. +// +// An adapter is expected to be hold all created workflow runs. +type Adapter interface { + // Init intializes the adapter + Init() error + + // Shutdown gracefully shuts down background jobs + Shutdown(ctx context.Context) error + + // SetupRepo ensures adapter connected to the repository. + // This usually includes adding repository watcher that does sparse-clone. + SetupRepo(ctx context.Context, repo syntax.ATURI) error + + // ListWorkflowDefs parses and returns all workflow definitions in the given + // repository at the specified revision + ListWorkflowDefs(ctx context.Context, repo syntax.ATURI, rev string) ([]WorkflowDef, error) + + // EvaluateEvent consumes a trigger event and returns a list of triggered + // workflow runs. It is expected to return immediately after scheduling the + // workflows. + EvaluateEvent(ctx context.Context, event Event) ([]WorkflowRun, error) + + // GetActiveWorkflowRun returns current state of specific workflow run. + // This method will be called regularly for active workflow runs. + GetActiveWorkflowRun(ctx context.Context, runId syntax.ATURI) (WorkflowRun, error) + + + + + // NOTE: baisically I'm not sure about this method. + // How to properly sync workflow.run states? + // + // for adapters with external engine, they will hold every past + // workflow.run objects. + // for adapters with internal engine, they... should also hold every + // past workflow.run objects..? + // + // problem: + // when spindle suffer downtime (spindle server shutdown), + // external `workflow.run`s might be unsynced in "running" or "pending" state + // same for internal `workflow.run`s. + // + // BUT, spindle itself is holding the runs, + // so it already knows unsynced workflows (=workflows not finished) + // therefore, it can just fetch them again. + // for adapters with internal engines, they will fail to fetch previous + // run. + // Leaving spindle to mark the run as "Lost" or "Failed". + // Because of _lacking_ adaters, spindle should be able to manually + // mark unknown runs with "lost" state. + // + // GetWorkflowRun : used to get background crawling + // XCodeCloud: ok + // Nixery: (will fail if unknown) -> spindle will mark workflow as failed anyways + // StreamWorkflowRun : used to notify real-time updates + // XCodeCloud: ok (but old events will be lost) + // Nixery: same. old events on spindle downtime will be lost + // + // + // To avoid this, each adapters should hold outbox buffer + // + // | + // v + + // StreamWorkflowRun(ctx context.Context) <-chan WorkflowRun + + + // ListActiveWorkflowRuns returns current list of active workflow runs. + // Runs where status is either Pending or Running + ListActiveWorkflowRuns(ctx context.Context) ([]WorkflowRun, error) + SubscribeWorkflowRun(ctx context.Context) <-chan WorkflowRun + + + + + // StreamWorkflowRunLogs streams logs for a running workflow execution + StreamWorkflowRunLogs(ctx context.Context, runId syntax.ATURI, handle func(line LogLine) error) error + + // CancelWorkflowRun attempts to stop a running workflow execution. + // It won't do anything when the workflow has already completed. + CancelWorkflowRun(ctx context.Context, runId syntax.ATURI) error +} diff --git a/spindle/models/pipeline2.go b/spindle/models/pipeline2.go new file mode 100644 index 00000000..fb800d46 --- /dev/null +++ b/spindle/models/pipeline2.go @@ -0,0 +1,124 @@ +package models + +import ( + "fmt" + "slices" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/tangled" +) + +// `sh.tangled.ci.event` +type Event struct { + SourceRepo syntax.ATURI // repository to find the workflow definition + SourceSha string // sha to find the workflow definition + TargetSha string // sha to run the workflow + // union type of: + // 1. PullRequestEvent + // 2. PushEvent + // 3. ManualEvent +} + +func (e *Event) AsRecord() tangled.CiEvent { + // var meta tangled.CiEvent_Meta + // return tangled.CiEvent{ + // Meta: &meta, + // } + panic("unimplemented") +} + +// `sh.tangled.ci.pipeline` +// +// Pipeline is basically a group of workflows triggered by single event. +type Pipeline2 struct { + Did syntax.DID + Rkey syntax.RecordKey + + Event Event // event that triggered the pipeline + WorkflowRuns []WorkflowRun // workflow runs inside this pipeline +} + +func (p *Pipeline2) AtUri() syntax.ATURI { + return syntax.ATURI(fmt.Sprintf("at://%s/%s/%s", p.Did, tangled.CiPipelineNSID, p.Rkey)) +} + +func (p *Pipeline2) AsRecord() tangled.CiPipeline { + event := p.Event.AsRecord() + runs := make([]string, len(p.WorkflowRuns)) + for i, run := range p.WorkflowRuns { + runs[i] = run.AtUri().String() + } + return tangled.CiPipeline{ + Event: &event, + WorkflowRuns: runs, + } +} + +// `sh.tangled.ci.workflow.run` +type WorkflowRun struct { + Did syntax.DID + Rkey syntax.RecordKey + + AdapterId string // adapter id + Name string // name of workflow run (not workflow definition name!) + Status WorkflowStatus // workflow status + // TODO: can add some custom fields like adapter-specific log-id +} + +func (r WorkflowRun) WithStatus(status WorkflowStatus) WorkflowRun { + return WorkflowRun{ + Did: r.Did, + Rkey: r.Rkey, + AdapterId: r.AdapterId, + Name: r.Name, + Status: status, + } +} + +func (r *WorkflowRun) AtUri() syntax.ATURI { + return syntax.ATURI(fmt.Sprintf("at://%s/%s/%s", r.Did, tangled.CiWorkflowRunNSID, r.Rkey)) +} + +func (r *WorkflowRun) AsRecord() tangled.CiWorkflowRun { + statusStr := string(r.Status) + return tangled.CiWorkflowRun{ + Adapter: r.AdapterId, + Name: r.Name, + Status: &statusStr, + } +} + +// `sh.tangled.ci.workflow.status` +type WorkflowStatus string + +var ( + WorkflowStatusPending WorkflowStatus = "pending" + WorkflowStatusRunning WorkflowStatus = "running" + WorkflowStatusFailed WorkflowStatus = "failed" + WorkflowStatusCancelled WorkflowStatus = "cancelled" + WorkflowStatusSuccess WorkflowStatus = "success" + WorkflowStatusTimeout WorkflowStatus = "timeout" + + activeStatuses [2]WorkflowStatus = [2]WorkflowStatus{ + WorkflowStatusPending, + WorkflowStatusRunning, + } +) + +func (s WorkflowStatus) IsActive() bool { + return slices.Contains(activeStatuses[:], s) +} + +func (s WorkflowStatus) IsFinish() bool { + return !s.IsActive() +} + +// `sh.tangled.ci.workflow.def` +// +// Brief information of the workflow definition. A workflow can be defined in +// any form. This is a common info struct for any workflow definitions +type WorkflowDef struct { + AdapterId string // adapter id + Name string // name or the workflow (usually the yml file name) + When any // events the workflow is listening to +} diff --git a/spindle/pipeline.go b/spindle/pipeline.go new file mode 100644 index 00000000..66e28946 --- /dev/null +++ b/spindle/pipeline.go @@ -0,0 +1,40 @@ +package spindle + +import ( + "context" + + "tangled.org/core/spindle/models" +) + +// createPipeline creates a pipeline from given event. +// It will call `EvaluateEvent` for all adapters, gather the triggered workflow +// runs, and constuct a pipeline record from them. pipeline record. It will +// return nil if no workflow run has triggered. +// +// NOTE: This method won't fail. If `adapter.EvaluateEvent` returns an error, +// the error will be logged but won't bubble-up. +// +// NOTE: Adapters might create sub-event on its own for workflows triggered by +// other workflow runs. +func (s *Spindle) createPipeline(ctx context.Context, event models.Event) (*models.Pipeline2) { + l := s.l + + pipeline := models.Pipeline2{ + Event: event, + } + + // TODO: run in parallel + for id, adapter := range s.adapters { + runs, err := adapter.EvaluateEvent(ctx, event) + if err != nil { + l.Error("failed to process trigger from adapter '%s': %w", id, err) + } + pipeline.WorkflowRuns = append(pipeline.WorkflowRuns, runs...) + } + + if len(pipeline.WorkflowRuns) == 0 { + return nil + } + + return &pipeline +} diff --git a/spindle/repomanager/repomanager.go b/spindle/repomanager/repomanager.go new file mode 100644 index 00000000..039fcb43 --- /dev/null +++ b/spindle/repomanager/repomanager.go @@ -0,0 +1,169 @@ +package repomanager + +import ( + "bufio" + "bytes" + "context" + "errors" + "fmt" + "os" + "os/exec" + "path/filepath" + "slices" + "strings" + + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/go-git/go-git/v5" + "github.com/go-git/go-git/v5/config" + "github.com/go-git/go-git/v5/plumbing/object" + kgit "tangled.org/core/knotserver/git" + "tangled.org/core/types" +) + +// RepoManager manages a `sh.tangled.repo` record with its git context. +// It can be used to efficiently fetch the filetree of the repository. +type RepoManager struct { + repoDir string + // TODO: it would be nice if RepoManager can be configured with different + // strategies: + // - use db as an only source for repo records + // - use atproto if record doesn't exist from the db + // - always use atproto + // hmm do we need `RepoStore` interface? + // now `DbRepoStore` and `AtprotoRepoStore` can implement both. + // all `RepoStore` objects will hold `KnotStore` interface, so they can + // source the knot store if needed. + + // but now we can't do complex queries like "get repo with issue count" + // that kind of queries will be done directly from `appview.DB` struct + // is graphql better tech for atproto? +} + +func New(repoDir string) RepoManager { + return RepoManager{ + repoDir: repoDir, + } +} + +// TODO: RepoManager can return file tree from repoAt & rev +// It will start syncing the repository if doesn't exist + +// RegisterRepo starts sparse-syncing repository with paths +func (m *RepoManager) RegisterRepo(ctx context.Context, repoAt syntax.ATURI, paths []string) error { + repoPath := m.repoPath(repoAt) + exist, err := isDir(repoPath) + if err != nil { + return fmt.Errorf("checking dir info: %w", err) + } + var sparsePaths []string + if !exist { + // init bare git repo + repo, err := git.PlainInit(repoPath, true) + if err != nil { + return fmt.Errorf("initializing repo: %w", err) + } + _, err = repo.CreateRemote(&config.RemoteConfig{ + Name: "origin", + URLs: []string{m.repoCloneUrl(repoAt)}, + }) + if err != nil { + return fmt.Errorf("configuring repo remote: %w", err) + } + } else { + // get sparse-checkout list + sparsePaths, err = func(path string) ([]string, error) { + var stdout bytes.Buffer + listCmd := exec.Command("git", "-C", path, "sparse-checkout", "list") + listCmd.Stdout = &stdout + if err := listCmd.Run(); err != nil { + return nil, err + } + + var sparseList []string + scanner := bufio.NewScanner(&stdout) + for scanner.Scan() { + line := strings.TrimSpace(scanner.Text()) + if line == "" { + continue + } + sparseList = append(sparseList, line) + } + if err := scanner.Err(); err != nil { + return nil, fmt.Errorf("scanning stdout: %w", err) + } + + return sparseList, nil + }(repoPath) + if err != nil { + return fmt.Errorf("parsing sparse-checkout list: %w", err) + } + + // add paths to sparse-checkout list + for _, path := range paths { + sparsePaths = append(sparsePaths, path) + } + sparsePaths = slices.Collect(slices.Values(sparsePaths)) + } + + // set sparse-checkout list + args := append([]string{"-C", repoPath, "sparse-checkout", "set", "--no-cone"}, sparsePaths...) + if err := exec.Command("git", args...).Run(); err != nil { + return fmt.Errorf("setting sparse-checkout list: %w", err) + } + return nil +} + +// SyncRepo sparse-fetch specific rev of the repo +func (m *RepoManager) SyncRepo(ctx context.Context, repo syntax.ATURI, rev string) error { + // TODO: fetch repo with rev. + panic("unimplemented") +} + +func (m *RepoManager) Open(repo syntax.ATURI, rev string) (*kgit.GitRepo, error) { + // TODO: don't depend on knot/git + return kgit.Open(m.repoPath(repo), rev) +} + +func (m *RepoManager) FileTree(ctx context.Context, repo syntax.ATURI, rev, path string) ([]types.NiceTree, error) { + if err := m.SyncRepo(ctx, repo, rev); err != nil { + return nil, fmt.Errorf("syncing git repo") + } + gr, err := m.Open(repo, rev) + if err != nil { + return nil, err + } + dir, err := gr.FileTree(ctx, path) + if err != nil { + if errors.Is(err, object.ErrDirectoryNotFound) { + return nil, nil + } + return nil, fmt.Errorf("loading file tree: %w", err) + } + return dir, err +} + +func (m *RepoManager) repoPath(repo syntax.ATURI) string { + return filepath.Join( + m.repoDir, + repo.Authority().String(), + repo.Collection().String(), + repo.RecordKey().String(), + ) +} + +func (m *RepoManager) repoCloneUrl(repo syntax.ATURI) string { + // 1. get repo & knot models from db. fetch it if doesn't exist + // 2. construct https clone url + panic("unimplemented") +} + +func isDir(path string) (bool, error) { + info, err := os.Stat(path) + if err == nil && info.IsDir() { + return true, nil + } + if os.IsNotExist(err) { + return false, nil + } + return false, err +} diff --git a/spindle/server.go b/spindle/server.go index caffbc72..3b0cbefc 100644 --- a/spindle/server.go +++ b/spindle/server.go @@ -49,6 +49,7 @@ type Spindle struct { l *slog.Logger n *notifier.Notifier engs map[string]models.Engine + adapters map[string]models.Adapter jq *queue.Queue cfg *config.Config ks *eventconsumer.Consumer -- 2.51.2