Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
64 changes: 64 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
name: ci

on:
push:
# master is covered by the run that release.yml calls before it tags, so pushing
# there does not need a second one
branches-ignore: [master]
workflow_call:

# A branch that is pushed again supersedes its own in-flight run, but master is left
# alone: its run gates the release.
concurrency:
group: ci-${{ github.workflow }}-${{ github.ref }}
cancel-in-progress: ${{ github.ref != 'refs/heads/master' }}

permissions:
contents: read

jobs:
lint:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v7

- uses: actions/setup-go@v7
with:
go-version-file: go.mod

- name: Tool versions
id: tools
run: |
echo "golangci=$(make -s print-GOLANGCI_VERSION)" >> "$GITHUB_OUTPUT"
echo "actionlint=$(make -s print-ACTIONLINT_VERSION)" >> "$GITHUB_OUTPUT"

- uses: golangci/golangci-lint-action@v9
with:
version: ${{ steps.tools.outputs.golangci }}

# the other half of `make lint`, so CI is never weaker than the local target
- name: Lint the workflows
env:
ACTIONLINT_VERSION: ${{ steps.tools.outputs.actionlint }}
run: |
go install "github.com/rhysd/actionlint/cmd/actionlint@${ACTIONLINT_VERSION}"
"$(go env GOPATH)/bin/actionlint"

- name: go.mod and go.sum are tidy
run: |
go mod tidy
git diff --exit-code go.mod go.sum

test:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v7

- uses: actions/setup-go@v7
with:
go-version-file: go.mod

- run: go build ./...

# -race because the resolver cache and the registry heartbeat are concurrent
- run: go test -race -count=1 ./...
68 changes: 68 additions & 0 deletions .github/workflows/release.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
name: release

on:
push:
branches: [master]

# Releases are serialised, never cancelled: two merges landing together must each get
# their turn at the tag rather than one clobbering the other's run.
concurrency:
group: release
cancel-in-progress: false

permissions:
contents: read

jobs:
# Nothing is tagged that has not passed the same lint and tests a PR does.
ci:
uses: ./.github/workflows/ci.yml

tag:
needs: ci
runs-on: ubuntu-latest
permissions:
contents: write
steps:
- uses: actions/checkout@v7
with:
fetch-depth: 0 # the whole tag history, to find the version to bump

- name: Work out the next version
id: version
run: |
set -euo pipefail

# A re-run of this workflow, or a commit tagged by hand, must not publish a
# second version of the same code.
if tag=$(git describe --exact-match --tags HEAD 2>/dev/null); then
echo "HEAD is already released as $tag, nothing to do"
exit 0
fi

latest=$(git tag -l 'v*' --sort=-v:refname | grep -E '^v[0-9]+\.[0-9]+\.[0-9]+$' | head -1 || true)
if [ -z "$latest" ]; then
next=v0.1.0
else
IFS=. read -r major minor patch <<<"${latest#v}"
next="v${major}.${minor}.$((patch + 1))"
fi

echo "next=$next" >> "$GITHUB_OUTPUT"
echo "Releasing $next (previous: ${latest:-none})"

- name: Tag and release
if: steps.version.outputs.next != ''
env:
GH_TOKEN: ${{ github.token }}
NEXT: ${{ steps.version.outputs.next }}
run: |
set -euo pipefail

git config user.name "github-actions[bot]"
git config user.email "41898282+github-actions[bot]@users.noreply.github.com"

git tag -a "$NEXT" -m "$NEXT"
git push origin "$NEXT"

gh release create "$NEXT" --title "$NEXT" --generate-notes
30 changes: 30 additions & 0 deletions .golangci.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
version: "2"

linters:
default: standard # errcheck, govet, ineffassign, staticcheck, unused
enable:
- bodyclose # a leaked response body in an HTTP middleware leaks a connection
- errorlint
- misspell
- wrapcheck # errors crossing a package boundary say which package they came from

settings:
wrapcheck:
extra-ignore-sigs:
# a middleware passing the next transport's error along is not the one that
# should be naming it
- .RoundTrip(

exclusions:
rules:
# tests close what they open, and report errors to t rather than to a caller who
# would need to know which package they came from
- path: _test\.go
linters:
- bodyclose
- errcheck
- wrapcheck

formatters:
enable:
- gofmt
20 changes: 20 additions & 0 deletions Makefile
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
GOLANGCI_VERSION ?= v2.13.2
ACTIONLINT_VERSION ?= v1.7.12

GOBIN ?= $(or $(shell go env GOBIN),$(shell go env GOPATH)/bin)

.PHONY: test lint install-tools

test:
go test -v ./...

lint:
$(GOBIN)/golangci-lint run ./...
$(GOBIN)/actionlint

install-tools:
go install github.com/golangci/golangci-lint/v2/cmd/golangci-lint@$(GOLANGCI_VERSION)
go install github.com/rhysd/actionlint/cmd/actionlint@$(ACTIONLINT_VERSION)

print-%:
@echo $($*)
50 changes: 50 additions & 0 deletions caller/caller.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
package caller

import (
"fmt"
"net/http"

"github.com/draincloud/callpack/caller/middleware"
)

type Caller struct {
client http.Client
mws []middleware.RoundTripperHandler
}

func New(client http.Client, mws ...middleware.RoundTripperHandler) *Caller {
return &Caller{
client: client,
mws: mws,
}
}

func (r *Caller) Use(middlewares ...middleware.RoundTripperHandler) {
r.mws = append(r.mws, middlewares...)
}

func (r *Caller) With(middlewares ...middleware.RoundTripperHandler) *Caller {
combined := make([]middleware.RoundTripperHandler, 0, len(r.mws)+len(middlewares))
combined = append(combined, r.mws...)
combined = append(combined, middlewares...)
return &Caller{client: r.client, mws: combined}
}

func (r *Caller) Client() *http.Client {
cl := r.client
if cl.Transport == nil {
cl.Transport = http.DefaultTransport
}
for _, handler := range r.mws {
cl.Transport = handler(cl.Transport)
}
return &cl
}

func (r *Caller) Do(req *http.Request) (*http.Response, error) {
resp, err := r.Client().Do(req)
if err != nil {
return nil, fmt.Errorf("caller: failed to execute request: %w", err)
}
return resp, nil
}
104 changes: 104 additions & 0 deletions caller/middleware/consul/consul.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,104 @@
package consul

import (
"context"
"fmt"
"net"
"net/http"
"strconv"
"sync"
"time"

"github.com/draincloud/callpack/caller/middleware"
"github.com/hashicorp/consul/api"
)

const DefaultTTL = time.Second

func Resolve(client *api.Client, ttl time.Duration) middleware.RoundTripperHandler {
if ttl <= 0 {
ttl = DefaultTTL
}
c := &cache{health: client.Health(), ttl: ttl, services: map[string]*service{}}

return func(next http.RoundTripper) http.RoundTripper {
return middleware.RoundTripperFunc(func(req *http.Request) (*http.Response, error) {
addr, err := c.address(req.Context(), req.URL.Hostname())
if err != nil {
return nil, err
}

out := req.Clone(req.Context())
if out.Host == "" {
out.Host = req.URL.Host
}
out.URL.Host = addr

return next.RoundTrip(out)
})
}
}

type cache struct {
health *api.Health
ttl time.Duration

mu sync.Mutex
services map[string]*service
}

type service struct {
mu sync.Mutex
addresses []string
next uint64
expires time.Time
}

func (c *cache) address(ctx context.Context, name string) (string, error) {
if name == "" {
return "", fmt.Errorf("consul: request has no host to resolve")
}

c.mu.Lock()
s, ok := c.services[name]
if !ok {
s = &service{}
c.services[name] = s
}
c.mu.Unlock()

s.mu.Lock()
defer s.mu.Unlock()

if time.Now().After(s.expires) {
addresses, err := c.lookup(ctx, name)
if err != nil {
return "", err
}
s.addresses, s.expires, s.next = addresses, time.Now().Add(c.ttl), 0
}

address := s.addresses[s.next%uint64(len(s.addresses))]
s.next++
return address, nil
}

func (c *cache) lookup(ctx context.Context, name string) ([]string, error) {
entries, _, err := c.health.Service(name, "", true, (&api.QueryOptions{}).WithContext(ctx))
if err != nil {
return nil, fmt.Errorf("consul: failed to resolve service %q: %w", name, err)
}
if len(entries) == 0 {
return nil, fmt.Errorf("consul: service %q has no healthy instances", name)
}

addresses := make([]string, 0, len(entries))
for _, entry := range entries {
host := entry.Service.Address
if host == "" {
host = entry.Node.Address
}
addresses = append(addresses, net.JoinHostPort(host, strconv.Itoa(entry.Service.Port)))
}
return addresses, nil
}
Loading