Vodiči

Go Crawling: Mastering Concurrent Data Fetching with Durability

Krenimo u indeksiranje: ovladavanje konkurentnim dohvaćanjem podataka uz trajnost

Produkcijski crawler nije petlja oko http.Get. To je trajni red poslova povezan s nepouzdanom mrežom: URL-ovi se umnožavaju, poslužitelji uzvraćaju pritisak, baze podataka privremeno postaju nedostupne, a uspješan odgovor može stići nekoliko trenutaka prije nego što stroj izgubi napajanje.

Ovaj vodič izrađuje ograničeni Go crawler s globalnim i po-autoritetskim ograničenjima konkurentnosti, otkazivanjem, ponovnim pokušajima za prolazne pogreške, deduplikacijom URL-ova, potvrđenim zakupima, oporavkom isteklih zakupa i trajnim SQLite frontierom. Njegov model isporuke je najmanje jednom: nakon pada sustava URL pod zakupom može se ponovno dohvatiti. Kontrolirane nuspojave stoga koriste idempotentne operacije baze podataka.

Preduvjeti i struktura projekta

Potrebni su vam Go 1.22 ili noviji, odlazni HTTP pristup i ovlaštenje za indeksiranje odabranih ciljeva. Implementacija koristi golang.org/x/net/html za raščlanjivanje HTML-a i modernc.org/sqlite, SQLite upravljački program u čistom Gou.

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

Arhitektura: zakupi prije HTTP-a

SQLite frontier istodobno je red i trajna kontrolna točka. Svaki je URL pending, leased ili done. Primarni ključ URL-a osigurava trajnu deduplikaciju nakon namjerno male politike kanonikalizacije: sheme i hostovi malim slovima, izostavljeni zadani portovi, prazne putanje pretvorene u / i uklonjeni fragmenti.

Preuzimanje je jedna atomska naredba baze podataka koja se potvrđuje prije nego što radnik čeka kapacitet ili obavi mrežni I/O. Svako uspješno preuzimanje dobiva svježi nasumični token. Obnova, oslobađanje i dovršavanje zahtijevaju upravo taj token. Ako se istekli zakup ponovno preuzme, zastarjeli radnik ne može produljiti ni dovršiti novu rezervaciju.

Otkucaj srca produljuje zakup tijekom čekanja semafora, ponovnih pokušaja i preuzimanja. Važno je da otkucaj srca ima vlastiti signal za zaustavljanje. Dovršavanje posla zaustavlja i pridružuje otkucaj srca prije otkazivanja konteksta posla, pa namjerno otkazivanje ne može pretvoriti obnovu u tijeku u lažan rezultat gubitka zakupa.

Infrastrukturni kvarovi baze podataka kvarovi su procesa, a ne signali praznog reda. Pogreške preuzimanja, brojanja frontiera, otkucaja srca i dovršavanja šire se do main, otkazuju ostale radnike i proizvode izlazni status različit od nule. To upravitelju usluge koji koristi Restart=on-failure omogućuje oporavak umjesto tihog prihvaćanja nedovršenih redaka.

Udaljeni GET i dalje se može ponoviti nakon pada sustava. Lokalna transakcija ne može učiniti udaljeni poslužitelj idempotentnim, stoga indeksirajte samo krajnje točke na kojima su prihvatljivi ponovljeni zahtjevi sigurnim metodama.

Broj radnika ograničava globalnu konkurentnost. Semafor označen kanonikalnim autoritetom ograničava svaki host. Preusmjeravanja su ograničena na isti autoritet kao prethodni zahtjev, pa svaki zahtjev u lancu preusmjeravanja ostaje pokriven već zadržanim semaforom. Ovo ograničenje na razini aplikacije također ograničava HTTP/2 tokove; sama ograničenja veza to ne čine.

Implementirajte crawler

Izradite main.go. Operacije baze podataka dobivaju zasebne rokove od dvije sekunde, kraće i od dvominutnog zakupa i od 25-sekundnog budžeta za zaustavljanje pri implementaciji. HTTP klijent neovisno ograničava uspostavljanje veze, TLS pregovaranje, zaglavlja odgovora i cijeli zahtjev.

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)
	}
}

Testirajte ograđivanje, ograničenja, preusmjeravanja i širenje pogrešaka

Izradite main_test.go. Ovi testovi pokrivaju ponovne pokušaje za prolazne pogreške, potiskivanje dvostrukih poveznica, konkurentnost na razini aplikacije, odbijanje preusmjeravanja, obradu isteklih zakupa, ograđivanje zastarjelih tokena, čisto zaustavljanje otkucaja srca i širenje pogrešaka baze podataka.

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")
	}
}

Pokrenite uobičajeno ponašanje, instrumentaciju konkurentnosti, statičku analizu i završnu izgradnju:

go test ./...
go test -race ./...
go vet ./...
go build -trimpath -o durable-crawler .

Pokrenite i promatrajte crawler

Svaki seed odobrava jedan točan kanonikalan autoritet, uključujući svaki nezadani port. Otkrivene poveznice izvan tih autoriteta odbacuju se. Preusmjeravanja su stroža: svaki skok mora zadržati autoritet prethodnog zahtjeva, čak i ako je drugi seed odobrio odredište.

./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

Strukturirani zapisnici prikazuju ponovne pokušaje, završne pogreške dohvaćanja, izgubljene zakupe, pogreške baze podataka i pogreške kontrolne točke. Expvar objavljuje brojače dovršavanja, ponovnih pokušaja, pogrešaka i oporavljenih zakupa. Njegova krajnja točka također izlaže informacije o izvođenju, stoga je držite privatnom i postavite upozorenja za trajni rast ponovnih pokušaja, ponovljeni gubitak zakupa, ponovna pokretanja usluge ili frontier koji prestane dovršavati poslove.

Performanse, sigurnost i pristojnost

Dodatni radnici pomažu samo kada se rad proteže preko dovoljno autoriteta. Velik broj radnika usmjeren na jednu stranicu uglavnom stvara čekatelje semafora. Zadržite -per-host konzervativnim, objavite smislen user agent i implementirajte robots politiku cilja prije indeksiranja sustava koje ne kontrolirate.

SQLite je prikladan za ograničene dokumentacijske stranice, revizije i operativne inventare. Distribuirani crawler treba poslužiteljsku bazu podataka, ali stroj stanja ostaje isti: atomski preuzmite s jedinstvenim tokenom, potvrdite prije HTTP-a, obnavljajte putem ograničenih upita, ogradite zastarjele vlasnike i zadržite iskren ugovor najmanje jednom.

Popis dopuštenih autoriteta nije potpuna SSRF zaštita. DNS odgovori mogu se promijeniti, odobrena imena mogu se razriješiti na privatnu infrastrukturu, a proxyji mijenjaju stvarnu mrežnu putanju. Provedite politiku izlaznog prometa kontroliranim proxyjem ili vatrozidom. Zadržite Goovu zadanu TLS provjeru; permisivne TLS postavke pretvaraju pogreške certifikata u sigurnosne propuste.

Tijelo odgovora je ograničeno, ali HTML raščlanjivanje i dalje stvara stablo dokumenta. Stranice s mnogo poveznica mogu dodijeliti znatnu količinu memorije. Za neprijateljske ili iznimno velike ulaze smanjite ograničenje tijela i nametnite najveći broj prihvaćenih poveznica po stranici.

Implementacija uz systemd

Ova jedinica hosta koristi dinamički identitet i dodjeljuje pristup za pisanje samo svojem direktoriju stanja kojim upravlja systemd. Prije spremanja kao durable-crawler.service pregledajte seed i putanju izvršne datoteke.

[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

Ovo su privilegirane naredbe hosta, a ne naredbe kontejnera. Prije instalacije pregledajte sve postojeće odredišne datoteke; install zamjenjuje datoteke na tim točnim putanjama.

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

Politiku izlaznog vatrozida primijenite zasebno, bez izlaganja loopback porta za metrike. Tijekom uvođenja testirajte DNS, TLS provjeru, preusmjeravanja i ponašanje proxyja jer se blokirana ovisnost često najprije pojavljuje kao običan istek vremena.

Uobičajeni produkcijski kvarovi

  • Dvostruki zahtjevi nakon ponovnog pokretanja: očekivani su kada proces umre prije potvrđivanja dovršavanja. Kontrolirane nizvodne učinke označite kanonikalnim URL-om ili drugim stabilnim idempotencijskim ključem.
  • Zastarjeli radnik dovršava ponovno preuzeti rad: tokeni preuzimanja ponovno su korišteni ili izostavljeni iz provjere vlasništva. Tretirajte token kao identitet ograđivanja rezervacije, a ne kao identifikator procesa.
  • Uspješan posao ostaje pod zakupom: zaustavljanje otkucaja srca otkazivanjem konteksta korištenog za obnovu stvara utrku otkazivanja. Koristite zaseban signal za zaustavljanje otkucaja srca, pridružite ga i tek tada otkažite kontekst posla.
  • Radnici nestaju dok retci ostaju: pogreška baze podataka tretirana je kao prazan red. Proširite infrastrukturne kvarove i izađite sa statusom različitim od nule ili ih ponovno pokušajte prema izričitoj ograničenoj politici.
  • Gašenje pokreće novi rad: ponovno provjerite otkazivanje odmah nakon svakog preuzimanja. Ako je gašenje stiglo tijekom poziva baze podataka, oslobodite tu točnu rezervaciju bez dohvaćanja.
  • Preusmjeravanja zaobilaze ograničenja autoriteta: broj veza ne ograničava HTTP/2 tokove. Zadržite aplikacijski semafor i odbijte ili zasebno regulirajte preusmjeravanja između autoriteta.
  • Ponovni pokušaji pojačavaju prekid rada: ograničite pokušaje, zadržite semafor kroz ponovne pokušaje, dodajte podrhtavanje i ograničite Retry-After.
  • SQLite razvija pritisak zaključavanja: držite transakcije bez mrežnog rada i izbjegavajte nepotrebne veze. Kontinuirano zauzeto višedijelno radno opterećenje preraslo je ovaj izbor pohrane.

Kontrolni popis za završnu provjeru

  1. Uspješno pokrenite go test -race ./....
  2. Potvrdite da globalna i po-autoritetska konkurentnost ostaju unutar konfiguracije.
  3. Provjerite da preusmjeravanja između autoriteta nikada ne kontaktiraju svoje odredište.
  4. Zaustavite otkucaj srca tijekom obnove i potvrdite da se posao u vlasništvu i dalje može dovršiti.
  5. Zatvorite bazu podataka ili je učinite nedostupnom i potvrdite da proces izlazi sa statusom različitim od nule.
  6. Pošaljite SIGTERM tijekom sporog zahtjeva i potvrdite da se rezervacija u vlasništvu oslobađa.
  7. Naglo prekinite proces, pustite da zakup istekne i potvrdite da se URL ponovno preuzima.
  8. Provjerite da zastarjeli tokeni ne mogu obnoviti, osloboditi ni dovršiti ponovno preuzeti rad.
  9. Testirajte dvostruke poveznice, prevelika tijela, prolazne statuse i završne statuse.
  10. Potvrdite da su direktorij stanja i slušatelj metrika dostupni samo predviđenim identitetima.
  11. Pregledajte kontrole izlaznog prometa i ovlaštenje za svaki cilj indeksiranja.

Trajnost nije samo zapisivanje URL-ova na disk. Ona znači očuvanje vlasništva kroz prekid, razlikovanje praznog reda od pokvarene infrastrukture i ograđivanje jučerašnjih radnika od današnjih preuzimanja. Kada su te granice izričite, konkurentnost prestaje biti kockanje i postaje operativni alat.

Portret autora bloga

Mihajlo

Ja sam Mihajlo — programer vođen znatiželjom, disciplinom i stalnom željom da stvorim nešto smisleno. Dijelim uvide, tutorijale i besplatne usluge kako bih pomogao drugima da pojednostave svoj rad i rastu u svijetu softvera i umjetne inteligencije koji se neprestano razvija.