Go Crawling: Mastering Concurrent Data Fetching with Durability
A production crawler is not a loop around http.Get. It is a durable work queue attached to an unreliable network: URLs multiply, hosts push back, databases become temporarily unavailable, and a successful response can arrive moments before the machine loses power.
This tutorial builds a bounded Go crawler with global and per-authority concurrency limits, cancellation, transient retries, URL deduplication, committed leases, expired-lease recovery, and a durable SQLite frontier. Its delivery model is at least once: after a crash, a leased URL may be fetched again. Controlled side effects therefore use idempotent database operations.
Prerequisites and project structure
You need Go 1.22 or newer, outbound HTTP access, and authorization to crawl the selected targets. The implementation uses golang.org/x/net/html for HTML parsing and modernc.org/sqlite, a pure-Go SQLite driver.
mkdir durable-crawler
cd durable-crawler
go mod init example.com/durable-crawler
go get golang.org/x/[email protected]
go get modernc.org/[email protected]
durable-crawler/
go.mod
go.sum
main.go
main_test.go
Architecture: leases before HTTP
The SQLite frontier is both the queue and the durable checkpoint. Every URL is pending, leased, or done. The URL primary key provides durable deduplication after a deliberately small canonicalization policy: lowercase schemes and hosts, omitted default ports, empty paths converted to /, and fragments removed.
Claiming is one atomic database statement that commits before the worker waits for capacity or performs network I/O. Every successful claim receives a fresh random token. Renewal, release, and completion require that exact token. If an expired lease is reclaimed, a stale worker cannot extend or complete the new reservation.
A heartbeat extends the lease during semaphore waits, retries, and downloads. Importantly, the heartbeat has its own stop signal. Finishing a job stops and joins the heartbeat before canceling the job context, so deliberate cancellation cannot turn an in-flight renewal into a false lease-loss result.
Database infrastructure failures are process failures, not empty-queue signals. Claim, frontier-count, heartbeat, and completion errors propagate to main, cancel the other workers, and produce a nonzero exit. That allows a service manager using Restart=on-failure to recover instead of silently accepting unfinished rows.
The remote GET can still be repeated after a crash. A local transaction cannot make a remote server idempotent, so crawl only endpoints where repeated safe-method requests are acceptable.
The worker count bounds global concurrency. A semaphore keyed by canonical authority bounds each host. Redirects are restricted to the same authority as the preceding request, so every request in a redirect chain remains covered by the semaphore already held. This application-level limit also bounds HTTP/2 streams; connection limits alone do not.
Implement the crawler
Create main.go. Database operations receive separate two-second deadlines, below both the two-minute lease and the 25-second deployment shutdown budget. The HTTP client independently bounds connection setup, TLS negotiation, response headers, and the entire request.
package main
import (
"bytes"
"context"
"crypto/rand"
"database/sql"
"encoding/hex"
"errors"
"expvar"
"flag"
"fmt"
"io"
"log/slog"
mrand "math/rand"
"net"
"net/http"
"net/url"
"os"
"os/signal"
"path/filepath"
"strconv"
"strings"
"sync"
"syscall"
"time"
"golang.org/x/net/html"
_ "modernc.org/sqlite"
)
var (
completed = expvar.NewInt("crawler_completed_total")
retries = expvar.NewInt("crawler_retries_total")
failures = expvar.NewInt("crawler_failures_total")
recovered = expvar.NewInt("crawler_expired_leases_total")
errCrossAuthority = errors.New("cross-authority redirect rejected")
errRedirectLimit = errors.New("redirect limit exceeded")
errLeaseLost = errors.New("lease ownership lost")
)
const (
dbTimeout = 2 * time.Second
leaseTime = 2 * time.Minute
)
type stringList []string
func (s *stringList) String() string { return strings.Join(*s, ",") }
func (s *stringList) Set(v string) error {
*s = append(*s, v)
return nil
}
type Frontier struct{ db *sql.DB }
func openFrontier(path string) (*Frontier, error) {
if err := os.MkdirAll(filepath.Dir(path), 0700); err != nil {
return nil, err
}
db, err := sql.Open("sqlite", path)
if err != nil {
return nil, err
}
db.SetMaxOpenConns(1)
db.SetMaxIdleConns(1)
for _, statement := range []string{
"PRAGMA journal_mode=WAL",
"PRAGMA synchronous=FULL",
"PRAGMA busy_timeout=2000",
`CREATE TABLE IF NOT EXISTS frontier (
url TEXT PRIMARY KEY,
status TEXT NOT NULL
CHECK(status IN ('pending','leased','done')),
lease_token TEXT,
lease_until INTEGER NOT NULL DEFAULT 0
)`,
`CREATE INDEX IF NOT EXISTS frontier_status_lease
ON frontier(status, lease_until)`,
} {
if _, err := db.Exec(statement); err != nil {
db.Close()
return nil, err
}
}
return &Frontier{db: db}, nil
}
func dbContext(parent context.Context) (context.Context, context.CancelFunc) {
return context.WithTimeout(parent, dbTimeout)
}
func claimToken() (string, error) {
var b [16]byte
if _, err := rand.Read(b[:]); err != nil {
return "", err
}
return hex.EncodeToString(b[:]), nil
}
func requireOwned(result sql.Result, err error) error {
if err != nil {
return err
}
n, err := result.RowsAffected()
if err != nil {
return err
}
if n != 1 {
return errLeaseLost
}
return nil
}
func (f *Frontier) enqueue(ctx context.Context, raw string) error {
ctx, cancel := dbContext(ctx)
defer cancel()
_, err := f.db.ExecContext(ctx,
`INSERT OR IGNORE INTO frontier(url,status)
VALUES(?, 'pending')`, raw)
return err
}
func (f *Frontier) recoverExpired(ctx context.Context) error {
ctx, cancel := dbContext(ctx)
defer cancel()
result, err := f.db.ExecContext(ctx, `UPDATE frontier
SET status='pending', lease_token=NULL, lease_until=0
WHERE status='leased' AND lease_until < ?`,
time.Now().UnixMilli())
if err == nil {
n, _ := result.RowsAffected()
recovered.Add(n)
}
return err
}
func (f *Frontier) claim(ctx context.Context) (string, string, bool, error) {
token, err := claimToken()
if err != nil {
return "", "", false, err
}
ctx, cancel := dbContext(ctx)
defer cancel()
now := time.Now().UnixMilli()
until := time.Now().Add(leaseTime).UnixMilli()
var raw string
err = f.db.QueryRowContext(ctx, `UPDATE frontier
SET status='leased', lease_token=?, lease_until=?
WHERE url = (
SELECT url FROM frontier
WHERE status='pending'
OR (status='leased' AND lease_until < ?)
ORDER BY url LIMIT 1
)
RETURNING url`, token, until, now).Scan(&raw)
if errors.Is(err, sql.ErrNoRows) {
return "", "", false, nil
}
if err != nil {
return "", "", false, err
}
return raw, token, true, nil
}
func (f *Frontier) renew(ctx context.Context, raw, token string) error {
ctx, cancel := dbContext(ctx)
defer cancel()
result, err := f.db.ExecContext(ctx, `UPDATE frontier
SET lease_until=?
WHERE url=? AND status='leased' AND lease_token=?`,
time.Now().Add(leaseTime).UnixMilli(), raw, token)
return requireOwned(result, err)
}
func (f *Frontier) release(raw, token string) error {
ctx, cancel := context.WithTimeout(context.Background(), dbTimeout)
defer cancel()
result, err := f.db.ExecContext(ctx, `UPDATE frontier
SET status='pending', lease_token=NULL, lease_until=0
WHERE url=? AND status='leased' AND lease_token=?`,
raw, token)
return requireOwned(result, err)
}
func (f *Frontier) finish(
ctx context.Context, raw, token string, links []string,
) error {
ctx, cancel := dbContext(ctx)
defer cancel()
tx, err := f.db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback()
for _, link := range links {
if _, err := tx.ExecContext(ctx,
`INSERT OR IGNORE INTO frontier(url,status)
VALUES(?, 'pending')`, link); err != nil {
return err
}
}
result, err := tx.ExecContext(ctx, `UPDATE frontier
SET status='done', lease_token=NULL, lease_until=0
WHERE url=? AND status='leased' AND lease_token=?`,
raw, token)
if err := requireOwned(result, err); err != nil {
return err
}
return tx.Commit()
}
func (f *Frontier) remaining(ctx context.Context) (int, error) {
ctx, cancel := dbContext(ctx)
defer cancel()
var n int
err := f.db.QueryRowContext(ctx,
`SELECT count(*) FROM frontier WHERE status!='done'`).Scan(&n)
return n, err
}
type hostLimits struct {
mu sync.Mutex
size int
hosts map[string]chan struct{}
}
func (h *hostLimits) forHost(authority string) chan struct{} {
h.mu.Lock()
defer h.mu.Unlock()
if ch := h.hosts[authority]; ch != nil {
return ch
}
ch := make(chan struct{}, h.size)
h.hosts[authority] = ch
return ch
}
type Crawler struct {
frontier *Frontier
client *http.Client
allowed map[string]struct{}
limits *hostLimits
logger *slog.Logger
maxBody int64
retryMax int
retryBase time.Duration
leaseEvery time.Duration
}
func canonical(raw string, base *url.URL) (string, *url.URL, error) {
u, err := url.Parse(raw)
if err != nil {
return "", nil, err
}
if base != nil {
u = base.ResolveReference(u)
}
u.Scheme = strings.ToLower(u.Scheme)
if u.Scheme != "http" && u.Scheme != "https" {
return "", nil, errors.New("unsupported URL scheme")
}
if u.User != nil || u.Hostname() == "" {
return "", nil, errors.New("credentials or empty host rejected")
}
host := strings.ToLower(u.Hostname())
port := u.Port()
if (u.Scheme == "http" && port == "80") ||
(u.Scheme == "https" && port == "443") {
port = ""
}
switch {
case port != "":
u.Host = net.JoinHostPort(host, port)
case strings.Contains(host, ":"):
u.Host = "[" + host + "]"
default:
u.Host = host
}
u.Fragment = ""
if u.Path == "" {
u.Path = "/"
}
return u.String(), u, nil
}
func redirectPolicy(req *http.Request, via []*http.Request) error {
if len(via) >= 10 {
return errRedirectLimit
}
if len(via) == 0 {
return nil
}
_, current, err := canonical(req.URL.String(), nil)
if err != nil {
return err
}
_, previous, err := canonical(via[len(via)-1].URL.String(), nil)
if err != nil {
return err
}
if current.Host != previous.Host {
return fmt.Errorf("%w: %s to %s",
errCrossAuthority, previous.Host, current.Host)
}
return nil
}
func extractLinks(body []byte, base *url.URL) []string {
doc, err := html.Parse(bytes.NewReader(body))
if err != nil {
return nil
}
var links []string
var walk func(*html.Node)
walk = func(n *html.Node) {
if n.Type == html.ElementNode && n.Data == "a" {
for _, a := range n.Attr {
if strings.EqualFold(a.Key, "href") {
if raw, _, err := canonical(a.Val, base); err == nil {
links = append(links, raw)
}
break
}
}
}
for child := n.FirstChild; child != nil; child = child.NextSibling {
walk(child)
}
}
walk(doc)
return links
}
func retryableStatus(status int) bool {
switch status {
case 408, 425, 429, 500, 502, 503, 504:
return true
}
return false
}
func retryAfter(value string) time.Duration {
if seconds, err := strconv.Atoi(strings.TrimSpace(value)); err == nil && seconds >= 0 {
return min(time.Duration(seconds)*time.Second, 30*time.Second)
}
if when, err := http.ParseTime(value); err == nil {
return min(max(time.Until(when), 0), 30*time.Second)
}
return 0
}
func (c *Crawler) fetch(
ctx context.Context, raw string,
) ([]string, bool, time.Duration, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, raw, nil)
if err != nil {
return nil, false, 0, err
}
req.Header.Set("User-Agent", "durable-crawler/1.0")
resp, err := c.client.Do(req)
if err != nil {
if ctx.Err() != nil {
return nil, false, 0, ctx.Err()
}
if errors.Is(err, errCrossAuthority) ||
errors.Is(err, errRedirectLimit) {
return nil, false, 0, err
}
return nil, true, 0, err
}
defer resp.Body.Close()
if retryableStatus(resp.StatusCode) {
_, _ = io.Copy(io.Discard, io.LimitReader(resp.Body, 4096))
return nil, true, retryAfter(resp.Header.Get("Retry-After")),
fmt.Errorf("retryable HTTP status %d", resp.StatusCode)
}
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return nil, false, 0,
fmt.Errorf("terminal HTTP status %d", resp.StatusCode)
}
if !strings.HasPrefix(
strings.ToLower(resp.Header.Get("Content-Type")), "text/html",
) {
return nil, false, 0, nil
}
body, err := io.ReadAll(io.LimitReader(resp.Body, c.maxBody+1))
if err != nil {
return nil, true, 0, err
}
if int64(len(body)) > c.maxBody {
return nil, false, 0,
fmt.Errorf("body exceeds %d bytes", c.maxBody)
}
return extractLinks(body, resp.Request.URL), false, 0, nil
}
func (c *Crawler) crawl(ctx context.Context, raw string) ([]string, error) {
for attempt := 0; ; attempt++ {
links, retry, delay, err := c.fetch(ctx, raw)
if err == nil || !retry || attempt >= c.retryMax {
return links, err
}
retries.Add(1)
if delay == 0 {
shift := min(attempt, 4)
base := c.retryBase * time.Duration(1<<shift)
delay = base +
time.Duration(mrand.Int63n(int64(base/2)+1))
}
c.logger.Warn("request retry",
"url", raw, "attempt", attempt+1,
"delay", delay, "error", err)
timer := time.NewTimer(delay)
select {
case <-ctx.Done():
timer.Stop()
return nil, ctx.Err()
case <-timer.C:
}
}
}
func (c *Crawler) maintainLease(
cancelJob context.CancelFunc, raw, token string,
) func() error {
stop := make(chan struct{})
done := make(chan error, 1)
interval := c.leaseEvery
if interval <= 0 {
interval = 30 * time.Second
}
go func() {
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-stop:
done <- nil
return
case <-ticker.C:
// Renewal is independent of deliberate job cancellation.
if err := c.frontier.renew(
context.Background(), raw, token,
); err != nil {
done <- err
cancelJob()
return
}
}
}
}()
return func() error {
close(stop)
return <-done
}
}
func reportFatal(
fatal chan<- error, cancel context.CancelFunc, err error,
) {
select {
case fatal <- err:
default:
}
cancel()
}
func (c *Crawler) reservationProblem(
raw, phase string, err error,
fatal chan<- error, cancel context.CancelFunc,
) bool {
if errors.Is(err, errLeaseLost) {
c.logger.Warn("lease ownership lost",
"url", raw, "phase", phase)
return false
}
wrapped := fmt.Errorf("%s %s: %w", phase, raw, err)
c.logger.Error("database operation failed",
"url", raw, "phase", phase, "error", err)
reportFatal(fatal, cancel, wrapped)
return true
}
func (c *Crawler) worker(
ctx context.Context, id int, fatal chan<- error,
cancelRun context.CancelFunc, wg *sync.WaitGroup,
) {
defer wg.Done()
for {
raw, token, ok, err := c.frontier.claim(ctx)
if err != nil {
if ctx.Err() == nil {
reportFatal(fatal, cancelRun,
fmt.Errorf("worker %d claim: %w", id, err))
}
return
}
if !ok {
n, err := c.frontier.remaining(ctx)
if err != nil {
if ctx.Err() == nil {
reportFatal(fatal, cancelRun,
fmt.Errorf("worker %d count frontier: %w", id, err))
}
return
}
if n == 0 {
return
}
timer := time.NewTimer(100 * time.Millisecond)
select {
case <-ctx.Done():
timer.Stop()
return
case <-timer.C:
continue
}
}
// Shutdown may have arrived while the blocking claim completed.
if ctx.Err() != nil {
if err := c.frontier.release(raw, token); err != nil {
if c.reservationProblem(
raw, "shutdown release", err, fatal, cancelRun,
) {
return
}
}
return
}
jobCtx, cancelJob := context.WithCancel(ctx)
stopLease := c.maintainLease(cancelJob, raw, token)
_, parsed, parseErr := canonical(raw, nil)
_, allowed := c.allowed[parsedHost(parsed)]
if parseErr != nil || !allowed {
leaseErr := stopLease()
cancelJob()
if leaseErr != nil {
if c.reservationProblem(
raw, "lease renewal", leaseErr, fatal, cancelRun,
) {
return
}
continue
}
if err := c.frontier.finish(
context.Background(), raw, token, nil,
); err != nil {
if c.reservationProblem(
raw, "completion", err, fatal, cancelRun,
) {
return
}
continue
}
completed.Add(1)
continue
}
limit := c.limits.forHost(parsed.Host)
select {
case limit <- struct{}{}:
case <-jobCtx.Done():
leaseErr := stopLease()
cancelJob()
if leaseErr != nil {
if c.reservationProblem(
raw, "lease renewal", leaseErr, fatal, cancelRun,
) {
return
}
}
if ctx.Err() != nil {
if err := c.frontier.release(raw, token); err != nil {
if c.reservationProblem(
raw, "shutdown release", err, fatal, cancelRun,
) {
return
}
}
return
}
continue
}
links, crawlErr := c.crawl(jobCtx, raw)
<-limit
// Join the heartbeat before canceling jobCtx. An in-flight
// renewal can finish normally instead of seeing context.Canceled.
leaseErr := stopLease()
cancelJob()
if leaseErr != nil {
if c.reservationProblem(
raw, "lease renewal", leaseErr, fatal, cancelRun,
) {
return
}
continue
}
if ctx.Err() != nil {
if err := c.frontier.release(raw, token); err != nil {
if c.reservationProblem(
raw, "shutdown release", err, fatal, cancelRun,
) {
return
}
}
return
}
var accepted []string
if crawlErr != nil {
failures.Add(1)
c.logger.Warn("crawl failed",
"worker", id, "url", raw, "error", crawlErr)
} else {
for _, link := range links {
canon, u, err := canonical(link, nil)
if err == nil {
if _, ok := c.allowed[u.Host]; ok {
accepted = append(accepted, canon)
}
}
}
}
if err := c.frontier.finish(
context.Background(), raw, token, accepted,
); err != nil {
if c.reservationProblem(
raw, "completion", err, fatal, cancelRun,
) {
return
}
continue
}
completed.Add(1)
}
}
func parsedHost(u *url.URL) string {
if u == nil {
return ""
}
return u.Host
}
func main() {
var seeds stringList
var workers, perHost, retryMax int
var dbPath, metricsAddr string
flag.Var(&seeds, "seed",
"seed URL; repeat to approve more authorities")
flag.IntVar(&workers, "workers", 16, "global worker limit")
flag.IntVar(&perHost, "per-host", 2,
"requests per authority")
flag.IntVar(&retryMax, "retries", 3,
"retries after the first attempt")
flag.StringVar(&dbPath, "db", "./state/frontier.db",
"SQLite frontier path")
flag.StringVar(&metricsAddr, "metrics", "127.0.0.1:9090",
"expvar address")
flag.Parse()
logger := slog.New(slog.NewJSONHandler(os.Stdout, nil))
if len(seeds) == 0 || workers < 1 ||
perHost < 1 || retryMax < 0 {
logger.Error("invalid configuration")
os.Exit(2)
}
frontier, err := openFrontier(dbPath)
if err != nil {
logger.Error("open frontier", "error", err)
os.Exit(1)
}
allowed := make(map[string]struct{})
for _, seed := range seeds {
raw, u, err := canonical(seed, nil)
if err != nil {
logger.Error("invalid seed",
"seed", seed, "error", err)
os.Exit(2)
}
allowed[u.Host] = struct{}{}
if err := frontier.enqueue(context.Background(), raw); err != nil {
logger.Error("enqueue seed", "error", err)
os.Exit(1)
}
}
if err := frontier.recoverExpired(context.Background()); err != nil {
logger.Error("recover leases", "error", err)
os.Exit(1)
}
transport := &http.Transport{
Proxy: http.ProxyFromEnvironment,
DialContext: (&net.Dialer{
Timeout: 3 * time.Second,
KeepAlive: 30 * time.Second,
}).DialContext,
ForceAttemptHTTP2: true,
MaxIdleConns: 100,
MaxIdleConnsPerHost: perHost,
MaxConnsPerHost: perHost,
IdleConnTimeout: 60 * time.Second,
TLSHandshakeTimeout: 5 * time.Second,
ResponseHeaderTimeout: 8 * time.Second,
}
client := &http.Client{
Transport: transport,
Timeout: 15 * time.Second,
CheckRedirect: redirectPolicy,
}
crawler := &Crawler{
frontier: frontier,
client: client,
allowed: allowed,
limits: &hostLimits{
size: perHost, hosts: make(map[string]chan struct{}),
},
logger: logger,
maxBody: 4 << 20,
retryMax: retryMax,
retryBase: 500 * time.Millisecond,
leaseEvery: 30 * time.Second,
}
signalCtx, stop := signal.NotifyContext(
context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
runCtx, cancelRun := context.WithCancel(signalCtx)
defer cancelRun()
fatal := make(chan error, 1)
mux := http.NewServeMux()
mux.Handle("/debug/vars", expvar.Handler())
metricsServer := &http.Server{
Addr: metricsAddr,
Handler: mux,
ReadHeaderTimeout: 2 * time.Second,
}
go func() {
if err := metricsServer.ListenAndServe();
!errors.Is(err, http.ErrServerClosed) {
reportFatal(fatal, cancelRun,
fmt.Errorf("metrics server: %w", err))
}
}()
var wg sync.WaitGroup
for id := 0; id < workers; id++ {
wg.Add(1)
go crawler.worker(runCtx, id, fatal, cancelRun, &wg)
}
wg.Wait()
var runErr error
select {
case runErr = <-fatal:
default:
}
checkpointCtx, checkpointCancel := context.WithTimeout(
context.Background(), dbTimeout)
if _, err := frontier.db.ExecContext(
checkpointCtx, "PRAGMA wal_checkpoint(TRUNCATE)",
); err != nil {
logger.Error("final WAL checkpoint", "error", err)
if runErr == nil {
runErr = fmt.Errorf("final WAL checkpoint: %w", err)
}
}
checkpointCancel()
shutdownCtx, shutdownCancel := context.WithTimeout(
context.Background(), 2*time.Second)
if err := metricsServer.Shutdown(shutdownCtx); err != nil && runErr == nil {
runErr = fmt.Errorf("metrics shutdown: %w", err)
}
shutdownCancel()
transport.CloseIdleConnections()
if err := frontier.db.Close(); err != nil && runErr == nil {
runErr = fmt.Errorf("close frontier: %w", err)
}
if runErr != nil {
logger.Error("crawler stopped with infrastructure failure",
"error", runErr)
os.Exit(1)
}
}
Test fencing, limits, redirects, and failure propagation
Create main_test.go. These tests cover transient retries, duplicate-link suppression, application-level concurrency, redirect rejection, expired-lease processing, stale-token fencing, clean heartbeat shutdown, and propagation of database failures.
package main
import (
"context"
"fmt"
"io"
"log/slog"
"net/http"
"net/http/httptest"
"path/filepath"
"sync"
"sync/atomic"
"testing"
"time"
)
func runWorkers(ctx context.Context, c *Crawler, count int) error {
runCtx, cancel := context.WithCancel(ctx)
defer cancel()
fatal := make(chan error, 1)
var wg sync.WaitGroup
for id := 0; id < count; id++ {
wg.Add(1)
go c.worker(runCtx, id, fatal, cancel, &wg)
}
wg.Wait()
select {
case err := <-fatal:
return err
default:
return nil
}
}
func TestCrawler(t *testing.T) {
var current, maximum, flaky, expiredHits, targetHits atomic.Int32
target := httptest.NewServer(http.HandlerFunc(
func(w http.ResponseWriter, r *http.Request) {
targetHits.Add(1)
io.WriteString(w, "must not be reached")
}))
defer target.Close()
server := httptest.NewServer(http.HandlerFunc(
func(w http.ResponseWriter, r *http.Request) {
n := current.Add(1)
defer current.Add(-1)
for {
old := maximum.Load()
if n <= old || maximum.CompareAndSwap(old, n) {
break
}
}
w.Header().Set("Content-Type", "text/html")
switch r.URL.Path {
case "/":
fmt.Fprintf(w,
`<a href="/a">A</a>`+
`<a href="/a">duplicate</a>`+
`<a href="/flaky">F</a>`+
`<a href="/jump">J</a>`)
case "/flaky":
if flaky.Add(1) <= 2 {
w.WriteHeader(http.StatusServiceUnavailable)
return
}
io.WriteString(w, "recovered")
case "/jump":
http.Redirect(w, r, target.URL, http.StatusFound)
case "/expired":
expiredHits.Add(1)
io.WriteString(w, "reclaimed")
default:
time.Sleep(10 * time.Millisecond)
io.WriteString(w, "done")
}
}))
defer server.Close()
frontier, err := openFrontier(
filepath.Join(t.TempDir(), "frontier.db"))
if err != nil {
t.Fatal(err)
}
defer frontier.db.Close()
root, parsed, err := canonical(server.URL, nil)
if err != nil {
t.Fatal(err)
}
if err := frontier.enqueue(context.Background(), root); err != nil {
t.Fatal(err)
}
client := server.Client()
client.CheckRedirect = redirectPolicy
crawler := &Crawler{
frontier: frontier,
client: client,
allowed: map[string]struct{}{parsed.Host: {}},
limits: &hostLimits{
size: 1, hosts: make(map[string]chan struct{}),
},
logger: slog.New(slog.NewTextHandler(io.Discard, nil)),
maxBody: 1 << 20,
retryMax: 3,
retryBase: time.Millisecond,
leaseEvery: 20 * time.Millisecond,
}
ctx, cancel := context.WithTimeout(
context.Background(), 3*time.Second)
defer cancel()
if err := runWorkers(ctx, crawler, 4); err != nil {
t.Fatal(err)
}
if maximum.Load() > 1 {
t.Fatalf("per-host limit violated: %d", maximum.Load())
}
if flaky.Load() != 3 {
t.Fatalf("expected three flaky requests, got %d", flaky.Load())
}
if targetHits.Load() != 0 {
t.Fatal("cross-authority redirect reached target")
}
var total, done int
if err := frontier.db.QueryRow(
"SELECT count(*) FROM frontier").Scan(&total); err != nil {
t.Fatal(err)
}
if err := frontier.db.QueryRow(
"SELECT count(*) FROM frontier WHERE status='done'",
).Scan(&done); err != nil {
t.Fatal(err)
}
if total != 4 || done != 4 {
t.Fatalf("unexpected frontier: total=%d done=%d", total, done)
}
expiredURL := server.URL + "/expired"
if _, err := frontier.db.Exec(`INSERT INTO frontier
(url,status,lease_token,lease_until)
VALUES(?, 'leased', 'dead-token', 0)`,
expiredURL); err != nil {
t.Fatal(err)
}
if err := frontier.recoverExpired(context.Background()); err != nil {
t.Fatal(err)
}
if err := runWorkers(ctx, crawler, 1); err != nil {
t.Fatal(err)
}
var status string
if err := frontier.db.QueryRow(
"SELECT status FROM frontier WHERE url=?", expiredURL,
).Scan(&status); err != nil {
t.Fatal(err)
}
if status != "done" || expiredHits.Load() != 1 {
t.Fatalf("reclaimed job not processed: status=%s hits=%d",
status, expiredHits.Load())
}
}
func TestClaimTokenFencesStaleWorker(t *testing.T) {
frontier, err := openFrontier(
filepath.Join(t.TempDir(), "frontier.db"))
if err != nil {
t.Fatal(err)
}
defer frontier.db.Close()
const job = "https://example.test/"
if err := frontier.enqueue(context.Background(), job); err != nil {
t.Fatal(err)
}
raw, stale, ok, err := frontier.claim(context.Background())
if err != nil || !ok {
t.Fatalf("first claim: ok=%v err=%v", ok, err)
}
if _, err := frontier.db.Exec(
"UPDATE frontier SET lease_until=0 WHERE url=?", raw,
); err != nil {
t.Fatal(err)
}
raw, current, ok, err := frontier.claim(context.Background())
if err != nil || !ok {
t.Fatalf("reclaim: ok=%v err=%v", ok, err)
}
if stale == current {
t.Fatal("reclaimed job reused its claim token")
}
if err := frontier.renew(context.Background(), raw, stale); err == nil {
t.Fatal("stale renewal succeeded")
}
if err := frontier.release(raw, stale); err == nil {
t.Fatal("stale release succeeded")
}
if err := frontier.finish(
context.Background(), raw, stale, nil,
); err == nil {
t.Fatal("stale completion succeeded")
}
if err := frontier.finish(
context.Background(), raw, current, nil,
); err != nil {
t.Fatalf("current owner could not complete: %v", err)
}
}
func TestHeartbeatStopsCleanly(t *testing.T) {
frontier, err := openFrontier(
filepath.Join(t.TempDir(), "frontier.db"))
if err != nil {
t.Fatal(err)
}
defer frontier.db.Close()
const job = "https://example.test/heartbeat"
if err := frontier.enqueue(context.Background(), job); err != nil {
t.Fatal(err)
}
raw, token, ok, err := frontier.claim(context.Background())
if err != nil || !ok {
t.Fatalf("claim: ok=%v err=%v", ok, err)
}
crawler := &Crawler{
frontier: frontier,
leaseEvery: time.Millisecond,
}
_, cancelJob := context.WithCancel(context.Background())
stopLease := crawler.maintainLease(cancelJob, raw, token)
time.Sleep(5 * time.Millisecond)
if err := stopLease(); err != nil {
t.Fatalf("deliberate heartbeat stop reported loss: %v", err)
}
cancelJob()
if err := frontier.finish(
context.Background(), raw, token, nil,
); err != nil {
t.Fatalf("completion after heartbeat stop failed: %v", err)
}
}
func TestWorkerReportsDatabaseFailure(t *testing.T) {
frontier, err := openFrontier(
filepath.Join(t.TempDir(), "frontier.db"))
if err != nil {
t.Fatal(err)
}
if err := frontier.db.Close(); err != nil {
t.Fatal(err)
}
crawler := &Crawler{
frontier: frontier,
logger: slog.New(slog.NewTextHandler(io.Discard, nil)),
}
ctx, cancel := context.WithTimeout(
context.Background(), time.Second)
defer cancel()
if err := runWorkers(ctx, crawler, 1); err == nil {
t.Fatal("database failure was mistaken for successful completion")
}
}
Exercise ordinary behavior, concurrency instrumentation, static analysis, and the final build:
go test ./...
go test -race ./...
go vet ./...
go build -trimpath -o durable-crawler .
Run and observe the crawler
Each seed approves one exact canonical authority, including any non-default port. Discovered links outside those authorities are discarded. Redirects are stricter: every hop must retain the preceding request’s authority, even if another seed approved the destination.
./durable-crawler \
-seed https://docs.example.org/ \
-workers 24 \
-per-host 2 \
-retries 3 \
-db ./state/frontier.db \
-metrics 127.0.0.1:9090
curl --fail --silent http://127.0.0.1:9090/debug/vars
Structured logs expose retries, terminal fetch failures, lost leases, database failures, and checkpoint failures. Expvar publishes completion, retry, failure, and recovered-lease counters. Its endpoint also exposes runtime information, so keep it private and alert on sustained retry growth, repeated lease loss, service restarts, or a frontier that stops completing.
Performance, security, and politeness
Additional workers help only when work spans enough authorities. A large worker count aimed at one site mostly creates semaphore waiters. Keep -per-host conservative, publish a meaningful user agent, and implement the target’s robots policy before crawling systems you do not control.
SQLite suits bounded documentation sites, audits, and operational inventories. A distributed crawler needs a server database, but the state machine remains the same: atomically claim with a unique token, commit before HTTP, renew through bounded queries, fence stale owners, and retain the honest at-least-once contract.
Authority allowlisting is not complete SSRF protection. DNS answers can change, approved names can resolve to private infrastructure, and proxies alter the effective network path. Enforce egress policy with a controlled proxy or firewall. Retain Go’s default TLS verification; permissive TLS settings turn certificate failures into security failures.
The response body is bounded, but HTML parsing still constructs a document tree. Link-dense pages can allocate substantial memory. For hostile or exceptionally large inputs, lower the body limit and impose a maximum number of admitted links per page.
Deploy with systemd
This host unit uses a dynamic identity and grants write access only to its systemd-managed state directory. Review the seed and executable path before saving it as durable-crawler.service.
[Unit]
Description=Durable bounded Go crawler
After=network-online.target
Wants=network-online.target
[Service]
Type=simple
DynamicUser=yes
StateDirectory=durable-crawler
ExecStart=/usr/local/bin/durable-crawler -seed https://docs.example.org/ -workers 24 -per-host 2 -retries 3 -db /var/lib/durable-crawler/frontier.db -metrics 127.0.0.1:9090
Restart=on-failure
RestartSec=5s
TimeoutStopSec=25s
NoNewPrivileges=yes
PrivateTmp=yes
ProtectSystem=strict
ProtectHome=yes
ProtectKernelTunables=yes
ProtectKernelModules=yes
ProtectControlGroups=yes
RestrictAddressFamilies=AF_INET AF_INET6
LockPersonality=yes
MemoryDenyWriteExecute=yes
[Install]
WantedBy=multi-user.target
These are privileged host commands, not container commands. Inspect any existing destination files before installation; install replaces files at those exact paths.
sudo stat /usr/local/bin/durable-crawler 2>/dev/null || true
sudo stat /etc/systemd/system/durable-crawler.service 2>/dev/null || true
sudo install -o root -g root -m 0755 \
./durable-crawler /usr/local/bin/durable-crawler
sudo install -o root -g root -m 0644 \
./durable-crawler.service \
/etc/systemd/system/durable-crawler.service
sudo systemctl daemon-reload
sudo systemctl enable --now durable-crawler.service
sudo systemctl status durable-crawler.service
sudo journalctl -u durable-crawler.service --since today
Apply outbound firewall policy separately, without exposing the loopback metrics port. During rollout, test DNS, TLS validation, redirects, and proxy behavior because a blocked dependency often first appears as an ordinary timeout.
Common production failures
- Duplicate requests after restart: expected when the process dies before committing completion. Key controlled downstream effects by canonical URL or another stable idempotency key.
- A stale worker completes reclaimed work: claim tokens were reused or omitted from an ownership check. Treat the token as the reservation’s fencing identity, not as a process identifier.
- A successful job remains leased: stopping a heartbeat by canceling the context used for renewal creates a cancellation race. Use a separate heartbeat stop signal, join it, and only then cancel the job context.
- Workers disappear while rows remain: a database error was treated like an empty queue. Propagate infrastructure failures and exit nonzero, or retry them under an explicit bounded policy.
- Shutdown starts fresh work: re-check cancellation immediately after every claim. If shutdown arrived during the database call, release that exact reservation without fetching.
- Redirects evade authority limits: connection counts do not bound HTTP/2 streams. Hold an application semaphore and reject or separately regulate cross-authority redirects.
- Retries amplify an outage: cap attempts, hold the semaphore across retries, add jitter, and bound
Retry-After. - SQLite develops lock pressure: keep transactions free of network work and avoid unnecessary connections. A continuously busy multi-process workload has outgrown this storage choice.
Final verification checklist
- Run
go test -race ./...successfully. - Confirm global and per-authority concurrency remain within configuration.
- Verify cross-authority redirects never contact their destination.
- Stop a heartbeat during renewal and confirm the owned job can still finish.
- Close or make the database unavailable and confirm the process exits nonzero.
- Send
SIGTERMduring a slow request and confirm the owned reservation is released. - Kill the process abruptly, let the lease expire, and confirm the URL is reclaimed.
- Verify stale tokens cannot renew, release, or finish reclaimed work.
- Test duplicate links, oversized bodies, transient statuses, and terminal statuses.
- Confirm the state directory and metrics listener are accessible only to intended identities.
- Review egress controls and authorization for every crawl target.
Durability is not merely writing URLs to disk. It is preserving ownership through interruption, distinguishing an empty queue from broken infrastructure, and fencing yesterday’s workers from today’s claims. Once those boundaries are explicit, concurrency stops being a gamble and becomes an operational tool.