package httpd import ( "encoding/json" "errors" "net/http" "strconv" "strings" "time" "github.com/barerepo/server/internal/artifact" "github.com/barerepo/server/internal/gitx" "github.com/barerepo/server/internal/proposal" "github.com/barerepo/server/internal/repo" "github.com/barerepo/server/internal/run" "github.com/barerepo/server/internal/store" "github.com/barerepo/server/internal/token" ) // pollWait holds a poll open, since the runner dials out and needs no inbound port. 15. const pollWait = 30 * time.Second // pollTick is how often a held poll looks for work. const pollTick = time.Second // runnerAuth reads the token and names its repository, since chapter 15 scopes one to each. func (s *Server) runnerAuth(r *http.Request, tok string) (*store.Token, error) { if tok == "" { return nil, errors.New("no token") } return s.DB.AccountForToken(r.Context(), token.Runner, tok) } type attachRequest struct { Token string `json:"token"` Hostname string `json:"hostname"` OS string `json:"os"` Arch string `json:"arch"` Labels []string `json:"labels"` } type attachResponse struct { RunnerID int64 `json:"runner_id"` PollInterval int `json:"poll_interval"` } func (s *Server) serveRunnerAttach(w http.ResponseWriter, r *http.Request) { var req attachRequest if err := json.NewDecoder(http.MaxBytesReader(w, r.Body, 1<<16)).Decode(&req); err != nil { http.Error(w, "malformed request", http.StatusBadRequest) return } t, err := s.runnerAuth(r, req.Token) if err != nil { http.Error(w, "that token is not valid", http.StatusUnauthorized) return } if req.Hostname == "" { http.Error(w, "a runner needs a hostname", http.StatusBadRequest) return } runner, err := s.DB.AttachRunner(r.Context(), t.ID, t.Scope, req.Hostname, req.OS, req.Arch, req.Labels) if err != nil { s.oops(w, r, err) return } writeJSON(w, attachResponse{RunnerID: runner.ID, PollInterval: int(pollWait.Seconds())}) } type jobResponse struct { JobID int64 `json:"job_id"` Repo string `json:"repo"` Ref string `json:"ref"` SHA string `json:"sha"` Command string `json:"command"` Image string `json:"image,omitempty"` CloneURL string `json:"clone_url"` JobToken string `json:"job_token"` } // serveRunnerPoll holds the request open until there is work or the wait ends. func (s *Server) serveRunnerPoll(w http.ResponseWriter, r *http.Request) { t, err := s.runnerAuth(r, r.URL.Query().Get("token")) if err != nil { http.Error(w, "that token is not valid", http.StatusUnauthorized) return } id, _ := strconv.ParseInt(r.URL.Query().Get("id"), 10, 64) // The id is the caller's to say, so it has to name a machine this token attached. 15. runner, err := s.runnerOf(r, t, id) if err != nil || runner.Repo != t.Scope { http.Error(w, "attach first", http.StatusNotFound) return } if err := s.DB.SeeRunner(r.Context(), runner.ID); err != nil { s.oops(w, r, err) return } deadline := time.After(pollWait) tick := time.NewTicker(pollTick) defer tick.Stop() for { job, err := s.DB.TakeJob(r.Context(), runner.Repo, *runner) if err != nil { s.oops(w, r, err) return } if job != nil { s.handOut(w, r, t, job) return } select { case <-r.Context().Done(): return case <-deadline: w.WriteHeader(http.StatusNoContent) return case <-tick.C: } } } // handOut gives a job to a runner, with a token scoped to that one job. func (s *Server) handOut(w http.ResponseWriter, r *http.Request, t *store.Token, job *store.Job) { // The long-lived token stays on the runner, and a job carries one that dies with it. jobToken, _, err := s.DB.CreateToken(r.Context(), token.Git, t.Account, job.Repo, jobLabel(job.ID)) if err != nil { s.oops(w, r, err) return } owner, name, _ := strings.Cut(job.Repo, "/") writeJSON(w, jobResponse{ JobID: job.ID, Repo: job.Repo, Ref: job.Ref, SHA: job.SHA, Command: job.Command, Image: job.Image, CloneURL: s.Cfg.Server.ExternalURL + "/" + owner + "/" + name, JobToken: jobToken, }) } type logRequest struct { Token string `json:"token"` JobID int64 `json:"job_id"` Seq int `json:"seq"` Chunk string `json:"chunk"` } func (s *Server) serveRunnerLog(w http.ResponseWriter, r *http.Request) { var req logRequest if err := json.NewDecoder(http.MaxBytesReader(w, r.Body, 1<<20)).Decode(&req); err != nil { http.Error(w, "malformed request", http.StatusBadRequest) return } _, runner, ok := s.jobRunner(w, r, req.Token, req.JobID) if !ok { return } if err := s.DB.AppendLog(r.Context(), req.JobID, runner.ID, req.Chunk); err != nil { http.Error(w, "that job is not running", http.StatusConflict) return } w.WriteHeader(http.StatusNoContent) } type doneRequest struct { Token string `json:"token"` JobID int64 `json:"job_id"` ExitCode int `json:"exit_code"` Duration int `json:"duration"` } // serveRunnerArtifact attaches one file to a release, which is chapter 22.5's job token permission. func (s *Server) serveRunnerArtifact(w http.ResponseWriter, r *http.Request) { tok := r.Header.Get("Barerepo-Token") jobID, _ := strconv.ParseInt(r.Header.Get("Barerepo-Job"), 10, 64) tag := r.Header.Get("Barerepo-Tag") file := r.Header.Get("Barerepo-File") job, ok := s.jobBearer(r, tok, jobID) if !ok { http.Error(w, "that token is not this job's", http.StatusUnauthorized) return } owner, name, ok := strings.Cut(job.Repo, "/") if !ok { http.Error(w, "that job has no repository", http.StatusConflict) return } // The token is scoped to one job and the job to one repository, so the tag must be in it. 22.5. dir, err := repo.Dir(s.Cfg.Paths.Repos, owner, name) if err != nil { s.oops(w, r, err) return } if _, err := gitx.ResolveRef(dir, "refs/tags/"+tag); err != nil { http.Error(w, "there is no release tagged "+tag, http.StatusNotFound) return } n, err := artifact.Put(s.Cfg.Paths.Artifacts, owner, name, tag, file, r.Body) if err != nil { http.Error(w, err.Error(), http.StatusBadRequest) return } if s.Log != nil { s.Log.Info("attached a release artifact", "repo", job.Repo, "tag", tag, "file", file, "bytes", n) } w.WriteHeader(http.StatusNoContent) } // jobBearer authorizes a job's own token, which chapter 22.5 scopes to one repository and one job. func (s *Server) jobBearer(r *http.Request, tok string, jobID int64) (*store.Job, bool) { if tok == "" || jobID == 0 { return nil, false } t, err := s.DB.AccountForToken(r.Context(), token.Git, tok) if err != nil { return nil, false } job, err := s.DB.Job(r.Context(), jobID) if err != nil || job == nil || job.Repo != t.Scope { return nil, false } // A requeued or finished job is over, and chapter 15 ends the token with it. if job.State != store.JobRunning { return nil, false } // The label is what the poll wrote, and it is what binds this token to this one job. if t.Label != jobLabel(jobID) { return nil, false } return job, true } // jobLabel names a job's token, in one place, so the poll and the check cannot drift apart. func jobLabel(id int64) string { return "job " + strconv.FormatInt(id, 10) } // dropJobToken revokes what the poll issued, since a job token outliving its job is a git credential nobody asked for, cloning forever. 15. func (s *Server) dropJobToken(r *http.Request, account string, jobID int64) { tokens, err := s.DB.TokensOf(r.Context(), account) if err != nil { s.log(r, err) return } label := jobLabel(jobID) for _, t := range tokens { if t.Kind != token.Git || t.Label != label { continue } if err := s.DB.DeleteToken(r.Context(), account, t.ID); err != nil { s.log(r, err) } } } // serveRunnerDone closes a job out and writes the result into the repository. func (s *Server) serveRunnerDone(w http.ResponseWriter, r *http.Request) { var req doneRequest if err := json.NewDecoder(http.MaxBytesReader(w, r.Body, 1<<16)).Decode(&req); err != nil { http.Error(w, "malformed request", http.StatusBadRequest) return } t, runner, ok := s.jobRunner(w, r, req.Token, req.JobID) if !ok { return } job, err := s.DB.FinishJob(r.Context(), req.JobID, runner.ID) if err != nil { http.Error(w, "that job is not running", http.StatusConflict) return } // Chapter 15: the job token expires when the job ends, and this is where a job ends. s.dropJobToken(r, t.Account, req.JobID) owner, name, _ := strings.Cut(job.Repo, "/") dir, err := repo.Dir(s.Cfg.Paths.Repos, owner, name) if err != nil { s.oops(w, r, err) return } rec := run.Record{ Runner: runner.Hostname, Labels: runner.Labels, Name: job.Name, Ref: job.Ref, Started: job.Started.Unix(), Duration: req.Duration, Exit: req.ExitCode, } if err := run.Append(r.Context(), dir, job.SHA, rec, job.Log); err != nil { s.oops(w, r, err) return } // Chapter 19.1 has run.failed and not run.succeeded, because a green build is not news. if req.ExitCode != 0 { s.note(r, store.Event{Kind: store.RunFailed, Actor: runner.Hostname, Repo: job.Repo, Ref: job.Ref, Number: proposal.Number(job.Ref), Title: job.Name, Detail: runner.Hostname + " ยท exit " + strconv.Itoa(req.ExitCode)}, 0) } w.WriteHeader(http.StatusNoContent) } // jobRunner checks that this token owns this job. func (s *Server) jobRunner(w http.ResponseWriter, r *http.Request, tok string, jobID int64) (*store.Token, *store.Runner, bool) { t, err := s.runnerAuth(r, tok) if err != nil { http.Error(w, "that token is not valid", http.StatusUnauthorized) return nil, nil, false } job, err := s.DB.Job(r.Context(), jobID) if err != nil || job.Repo != t.Scope { http.NotFound(w, r) return nil, nil, false } // The machine holding the job has to be this token's, or the repository's other runners could write into a build they are not running. 15. runner, err := s.runnerOf(r, t, job.RunnerID) if err != nil { http.NotFound(w, r) return nil, nil, false } if err := s.DB.SeeRunner(r.Context(), runner.ID); err != nil { s.log(r, err) } return t, runner, true } // runnerOf finds one machine among those that attached with this token, and no others. func (s *Server) runnerOf(r *http.Request, t *store.Token, id int64) (*store.Runner, error) { attached, err := s.DB.RunnersByToken(r.Context(), t.Account) if err != nil { return nil, err } mine := attached[t.ID] for i := range mine { if mine[i].ID == id { return &mine[i], nil } } return nil, store.ErrNotFound } func writeJSON(w http.ResponseWriter, v any) { w.Header().Set("Content-Type", "application/json") json.NewEncoder(w).Encode(v) }