Tutorials

Go Crawling: Mastering Concurrent Data Fetching with Durability

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

  1. Run go test -race ./... successfully.
  2. Confirm global and per-authority concurrency remain within configuration.
  3. Verify cross-authority redirects never contact their destination.
  4. Stop a heartbeat during renewal and confirm the owned job can still finish.
  5. Close or make the database unavailable and confirm the process exits nonzero.
  6. Send SIGTERM during a slow request and confirm the owned reservation is released.
  7. Kill the process abruptly, let the lease expire, and confirm the URL is reclaimed.
  8. Verify stale tokens cannot renew, release, or finish reclaimed work.
  9. Test duplicate links, oversized bodies, transient statuses, and terminal statuses.
  10. Confirm the state directory and metrics listener are accessible only to intended identities.
  11. 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.

Blog author portrait

Mihajlo

I’m Mihajlo — a developer driven by curiosity, discipline, and the constant urge to create something meaningful. I share insights, tutorials, and free services to help others simplify their work and grow in the ever-evolving world of software and AI.