crawler

package module
v1.0.3 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 4, 2026 License: MIT Imports: 32 Imported by: 0

README

crawler (hermes-crawler)

A high-throughput, restart-safe, politeness-enforcing web crawler engine and microservice in Go.

CI Pipeline Go Reference

crawler (Hermes Crawler) is a production-grade web crawling engine and microservice implemented in Go (github.com/hmza-hb/crawler). It accepts seed URLs or XML sitemaps, traverses hypermedia links under atomic resource budgets, and persists documents into PostgreSQL with HTTP ETag and Last-Modified conditional re-validation.

It operates both as an embeddable Go package and as a standalone REST API microservice daemon (cmd/crawld).

                       ┌─────────────────────────────────────────┐
                       │           crawld REST API               │
                       └──────────────────┬──────────────────────┘
                                          │
                                          ▼
                       ┌───────────────────────────────────────────┐
                       │    Crawler Worker Pool (crawl.go)         │
                       └─────┬───────────────────────────────┬─────┘
                             │                               │
                             ▼                               ▼
        ┌───────────────────────────────┐     ┌───────────────────────────────┐
        │     MemoryFrontier Queue      │     │      PgFrontier (SKIP LOCKED) │
        └──────────────┬────────────────┘     └──────────────┬────────────────┘
                       │                                     │
                       └──────────────────┬──────────────────┘
                                          │
                                          ▼
                      ┌───────────────────────────────────────┐
                      │    Polite HTTP Fetcher (fetch.go)     │
                      └──────────────────┬────────────────────┘
                                         │
              ┌──────────────────────────┴──────────────────────────┐
              ▼                                                     ▼
 ┌─────────────────────────┐                               ┌─────────────────┐
 │ RFC 9309 Robots Parser  │                               │  SSRF IP Guard  │
 └─────────────────────────┘                               └─────────────────┘

Why crawler?

Crawling web resources at scale while maintaining system safety and legal/politeness compliance presents non-trivial distributed engineering challenges:

  1. Strict Politeness Compliance: Uncontrolled crawlers risk overloading target hosts or violating site access policies. crawler incorporates a dedicated RFC 9309 compliant robots.txt parser and domain-keyed token bucket rate limiting.
  2. Re-validation & Bandwidth Efficiency: Downloading identical web pages repeatedly wastes network I/O and storage. crawler executes conditional HTTP GET requests (If-None-Match / If-Modified-Since), processing 304 Not Modified responses to prevent write amplification.
  3. Restart Safety & Non-Blocking Leasing: Node crashes during long-running crawls can corrupt queue state. crawler provides a PostgreSQL frontier (PgFrontier) utilizing FOR UPDATE SKIP LOCKED for concurrent work leasing and automatic lease recovery without external message brokers.
  4. Zero-Trust Network Isolation: Malicious links or open redirects can trigger Server-Side Request Forgery (SSRF) against cloud infrastructure endpoints (169.254.169.254, internal services). crawler intercepts DNS resolutions and redirect chains, validating target IP addresses against private CIDR ranges prior to HTTP transport.

Engineering Guarantees & Features

  • RFC 9309 Compliant Engine: Parses wildcards (*, $), Crawl-delay, group precedence, treats 404 Not Found as permissive, and fails closed on network errors.
  • Conditional GET Re-validation: Tracks HTTP ETag and Last-Modified response metadata, avoiding body downloads when content remains unchanged.
  • SSRF Isolation & IP Range Guards: Re-evaluates every DNS lookup and redirect hop against private ranges (10.0.0.0/8, 172.16.0.0/12, 192.168.0.0/16, 169.254.169.254, loopback).
  • Distributed Work Leasing (SKIP LOCKED): Multi-node URL scheduling backed by PostgreSQL row-level locks without lock contention.
  • Atomic Budget Enforcement: Hard bounds across total pages, retained bytes, wall-clock duration, crawl depth, and per-host page limits.
  • Dual Infrastructure Modes: Embeddable Go package and standalone REST API daemon (cmd/crawld).
  • Prometheus Observability: Native exporter tracking fetch latencies, queue depth, status codes, and HTTP server metrics.

System Architecture & Domain Subsystems

The core architecture isolates concerns into four primary component abstractions:

┌─────────────────────────────────────────────────────────────────────────────────┐
│                                   Crawler                                       │
│    Worker Pool Coordinator · Link Canonicalization · Budget Enforcement         │
└───────┬─────────────────────────────────────────────────────────────────┬───────┘
        │                                                                 │
        ▼                                                                 ▼
┌───────────────────────────────┐                 ┌───────────────────────────────┐
│           Frontier            │                 │            Fetcher            │
│ Memory / PostgreSQL Queueing  │                 │ Polite HTTP Client / SSRF     │
└───────────────┬───────────────┘                 └───────────────┬───────────────┘
                │                                                 │
                └────────────────────────┬────────────────────────┘
                                         ▼
                         ┌───────────────────────────────┐
                         │             Store             │
                         │ Document & Metadata Storage   │
                         └───────────────────────────────┘
Subsystem Breakdown
  1. Fetcher (fetch.go): Manages HTTP connection pooling, domain-keyed token bucket rate limiting (internal/platform/ratelimit), RFC 9309 parsing, and conditional GET re-validation.
  2. Frontier (frontier.go, pgfrontier.go): Manages URL traversal state (pending -> leased -> completed/failed). MemoryFrontier provides an in-memory priority queue, while PgFrontier handles PostgreSQL persistence with row locking.
  3. Store (store.go, pgstore.go): Persistence layer storing HTTP status codes, headers, execution latency, raw/normalized SHA-256 content hashes, and extracted hypermedia links.
  4. Crawler (crawl.go): Multi-threaded worker pool coordinator executing depth traversal, link canonicalization, HTML/Sitemap parsing, and budget tracking.

→ Detailed Architecture Document


Deep-Dive Engineering Decisions

1. PostgreSQL SKIP LOCKED vs External Message Queues

Rather than coupling the system to complex external message queues (e.g., Kafka, RabbitMQ, or Redis Streams), PgFrontier uses PostgreSQL FOR UPDATE SKIP LOCKED. Worker nodes query pending URLs concurrently without waiting for locked rows. This eliminates lock contention while reducing operational setup overhead to PostgreSQL.

2. Lease Expiry & Fault Recovery

When a worker claims a URL from PgFrontier, a lease timestamp (leased_until) is set. If a crawler worker process panics or dies mid-crawl, the lease automatically expires. Subsequent worker polls detect and reclaim expired leases seamlessly without manual intervention.

3. Domain-Keyed Token Bucket Rate Limiting

Rate limiting operates at the registrable domain level (via publicsuffix) rather than globally or per-IP. This ensures domain politeness (PerHostDelay) while allowing maximum parallel throughput across distinct web hosts.

4. Zero-Trust SSRF Interception

To prevent Server-Side Request Forgery (SSRF), DNS resolution is performed explicitly prior to HTTP transport execution. Target IP addresses are evaluated against reserved private CIDRs, loopback addresses, and cloud provider metadata IPs (169.254.169.254). Redirect chains are re-checked at every hop.


Quick Start & Usage

1. Embeddable Go Library

Add the package to your go.mod:

go get github.com/hmza-hb/crawler

Initialize and run a bounded crawl:

package main

import (
	"context"
	"fmt"
	"time"

	"github.com/hmza-hb/crawler"
)

func main() {
	cfg := crawler.DefaultConfig()
	store := crawler.NewMemoryStore()
	frontier := crawler.NewMemoryFrontier(cfg)

	fetcher, err := crawler.NewHTTPFetcher(cfg, crawler.Options{Store: store})
	if err != nil {
		panic(err)
	}

	loop, err := crawler.New(fetcher, frontier, store, cfg, crawler.Options{})
	if err != nil {
		panic(err)
	}

	res, err := loop.Crawl(context.Background(), crawler.CrawlOptions{
		Seeds:       []string{"https://example.com"},
		FollowLinks: true,
		Budget: crawler.Budget{
			MaxPages:    50,
			MaxDuration: 2 * time.Minute,
		},
	})
	if err != nil && err != crawler.ErrCrawlStopped {
		panic(err)
	}

	fmt.Printf("Fetched: %d, Not Modified: %d, Stopped By: %s\n",
		res.Fetched, res.NotModified, res.Stats.StoppedBy)
}
2. Standalone Microservice (crawld)
Running with Docker Compose
make docker-up
Running locally via Go CLI
export DATABASE_URL='postgres://postgres:postgres@localhost:5432/crawler?sslmode=disable'
export HTTP_ADDR=':8080'

go run ./cmd/crawld

HTTP REST API Endpoints

Method Endpoint Description
POST /v1/crawl Trigger a bounded, multi-page crawl run
POST /v1/fetch Politely fetch a single URL with persistence
GET /v1/documents/{id} Retrieve stored document metadata by UUID
GET /v1/documents?url={url} Query the last stored document for a target URL
GET /v1/hosts/{host}/robots Inspect cached robots.txt rules and crawl delay
GET /v1/stats?run_id={id} Retrieve queue depth and run statistics
GET /healthz, /readyz Liveness and database connectivity readiness probes
GET /metrics Prometheus metrics endpoint
Example REST Request & Response
Single URL Fetch Request
curl -X POST http://localhost:8080/v1/fetch \
  -H 'Content-Type: application/json' \
  -d '{"url": "https://example.com"}'
Response Payload
{
  "outcome": "fetched",
  "url": "https://example.com/",
  "status": 200,
  "status_text": "OK",
  "content_type": "text/html; charset=UTF-8",
  "content_hash": "a67f3...e3b1",
  "etag": "\"314b5b701509a25b3992070e6508d519\"",
  "last_modified": "2024-01-30T10:00:00Z",
  "not_modified": false,
  "filtered": false,
  "body_truncated": false,
  "bytes": 1256,
  "elapsed_ms": 142
}

Configuration Reference

The service and library are configured via Config or environment variables:

Environment Variable Description Default
DATABASE_URL PostgreSQL connection DSN Required for crawld
HTTP_ADDR Listen address for API microservice :8080
CRAWLER_USER_AGENT HTTP User-Agent header sent upstream HermesCrawler/1.0 (+https://github.com/hmza-hb/crawler)
CRAWLER_CONCURRENCY Concurrent worker goroutines per run 10
CRAWLER_PER_HOST_DELAY Minimum delay between requests to same host 250ms
CRAWLER_MAX_BYTES Maximum allowed HTTP response body size 10485760 (10MB)
CRAWLER_MAX_RETRIES Retry attempts for transient failures 2
CRAWLER_ALLOW_PRIVATE_HOSTS Override SSRF guard to permit loopback/private IPs false

Observability & Health Monitoring

1. Probes
  • GET /healthz — Service liveness probe returning HTTP 200.
  • GET /readyz — Database readiness probe executing a ping against PostgreSQL.
2. Prometheus Metrics

GET /metrics exposes standardized operational metrics:

  • crawler_fetch_duration_seconds: Histogram measuring HTTP fetch latency.
  • crawler_fetched_documents_total: Counter by HTTP status code and outcome.
  • crawler_frontier_depth: Gauge of current pending URLs in queue.

Development & Verification

Running Unit & Integration Tests
make test
Running Tests with Race Detection
go test -v -race ./...
Coverage Reports
make test-coverage
Performance Benchmarking

Execute Go micro-benchmarking suite for robots parsing, URL normalization, and frontier queueing:

go test -bench=. -benchmem ./...

Project Structure

.
├── cmd/
│   └── crawld/                 # REST API microservice entrypoint & CLI flags
├── internal/
│   └── platform/               # Platform infrastructure (db pool, httpx server, observe, ratelimit)
├── robotstxt/                  # Custom RFC 9309 robots.txt parser engine
├── docs/
│   ├── architecture.md         # Deep-dive system design document
│   └── openapi.yaml            # OpenAPI 3.0 REST API specification
├── migrations/                 # Embedded PostgreSQL migrations (`001_init.sql`, etc.)
├── docker-compose.yml          # Local development stack (Postgres + Service)
├── Dockerfile                  # Multi-stage production container build
├── Makefile                    # Build, test, and container targets
├── config.go                   # Core configuration structures & validation
├── crawl.go                    # Crawler worker loop & execution manager
├── fetch.go                    # Polite HTTP fetcher & security guards
├── frontier.go                 # In-memory frontier queue implementation
├── pgfrontier.go               # PostgreSQL SKIP LOCKED frontier implementation
├── store.go                    # Storage interfaces & memory storage
└── pgstore.go                  # PostgreSQL document & run persistence

Limitations & Scope

  • No JavaScript Execution: crawler parses static HTML and XML sitemaps. Single-page applications (SPAs) requiring client-side JS rendering are not executed.
  • Single-Process Memory Mode: MemoryFrontier operates in-memory for standalone single-process usage. Multi-node distributed deployment requires PgFrontier backed by PostgreSQL.

Roadmap

  • Dynamic adaptive per-host crawl delay based on server response latency.
  • Pluggable Object Storage (AWS S3 / GCS) backend for raw HTML content bodies.
  • JSON-LD and Schema.org metadata extraction pipeline.

Contributing

Contributions are welcome! Please follow these guidelines:

  1. Open an issue describing the proposed bug fix or feature.
  2. Ensure all changes include unit tests covering new behavior.
  3. Verify that make test and make lint pass before submitting pull requests.

License

MIT © 2026 Hamza & Contributors

Documentation

Overview

Package crawler is a polite, resumable, horizontally scalable web crawler.

It owns the canonical Document type — "a page we fetched, with its bytes, its HTTP metadata, and the content hash everything else keys off" — plus the frontier, robots compliance, politeness, conditional re-fetch, and the content-type and size policy that keep a crawl cheap.

The module has two halves that can be used independently:

  • Fetcher: fetch one URL, correctly, once. No queue, no state.
  • Crawler: a frontier-driven loop over many URLs with budgets.

cmd/crawld exposes both over HTTP so the crawler can be sold as a service without the rest of the platform.

Index

Constants

View Source
const DefaultFrontierSize = 200_000

DefaultFrontierSize bounds the in-memory frontier so a crawl that follows too many links fails loudly instead of exhausting the machine.

View Source
const DefaultMaxLinks = 500

DefaultMaxLinks bounds link extraction so a navigation-heavy page cannot dominate a crawl.

Variables

View Source
var (
	// ErrNotFound is returned when a document ID is unknown.
	ErrNotFound = errors.New("crawler: document not found")
	// ErrStoreClosed is returned once the store has been closed.
	ErrStoreClosed = errors.New("crawler: store is closed")
)

Store errors.

View Source
var ErrCrawlStopped = errors.New("crawler: crawl stopped")

ErrCrawlStopped is returned by Crawl when a budget, the context, or a fatal dependency error ends the run. A stopped crawl is not a failed crawl: everything fetched before the stop is valid and already stored.

View Source
var ErrFrontierClosed = errors.New("crawler: frontier is closed")

ErrFrontierClosed is returned once the frontier has been closed.

View Source
var Migrations embed.FS

Migrations holds the crawler's schema. Every module owns its own migrations and the lead-engine CLI passes them all to db.Migrate together, so a new module's schema ships with that module rather than in a central place.

Functions

func CharsetOf

func CharsetOf(header string) string

CharsetOf extracts the charset parameter from a Content-Type header.

func ContentTypeOf

func ContentTypeOf(header string) string

ContentTypeOf splits a media type from its parameters, lowercased.

func DefaultContentTypes

func DefaultContentTypes() []string

DefaultContentTypes are the media types this crawler will store and parse. Everything else is fetched for its status and then discarded, which is how a crawl stays cheap on sites that serve a hundred megabytes of media.

func IsHTMLContentType

func IsHTMLContentType(ct string) bool

IsHTMLContentType reports whether a media type is worth parsing for links and facts.

func IsRetryableStatus

func IsRetryableStatus(code int) bool

IsRetryableStatus reports whether a status justifies an immediate retry. 4xx other than 429 will not change by asking again.

func IsSameSite

func IsSameSite(a, b *url.URL) bool

IsSameSite reports whether two URLs are plausibly the same site: the same host, or hosts sharing a registrable domain. It is the test a caller should use before treating two discovered URLs as two companies.

func IsXMLContentType

func IsXMLContentType(ct string) bool

IsXMLContentType reports whether a media type is a sitemap or feed.

func MigrationSource

func MigrationSource() db.Source

MigrationSource returns the crawler's migrations for db.Migrate.

func NewHTTPClient

func NewHTTPClient(cfg Config) *http.Client

NewHTTPClient builds the HTTP client the crawler uses: bounded timeouts, compression, a connection pool sized for concurrency, and a transport that refuses to dial private addresses unless explicitly allowed.

func NormalizeURL

func NormalizeURL(u *url.URL) (string, error)

NormalizeURL canonicalises a parsed URL: lowercased scheme and host, default port removed, fragment dropped, dot segments resolved. The query is preserved because for some sites it selects the content.

Exported because other modules dedupe against the crawler's keys and must produce byte-identical strings.

func ParseSitemap

func ParseSitemap(r io.Reader) (entries []SitemapEntry, index []SitemapIndexEntry, err error)

ParseSitemap reads either a urlset or a sitemapindex. Anything it cannot parse returns an error rather than a partial result, so a caller can distinguish "not a sitemap" from "a sitemap with three entries".

func PublicURL

func PublicURL(raw string) (*url.URL, error)

PublicURL validates that a URL discovered from an untrusted source is one the crawler is willing to contact, and returns it normalized.

Discovered URLs arrive from search results, sitemaps, certificate logs and third-party pages, none of which are trustworthy: any of them can name http://169.254.169.254/, http://localhost:5432/ or file:///etc/passwd. Every discovered URL must pass through this function before a caller acts on it.

It rejects non-http(s) schemes, embedded credentials, and hosts inside private, loopback, link-local, CGNAT and other unroutable ranges. A hostname that is not an IP literal passes here and is re-checked after DNS resolution by the crawler's transport, so this is the cheap first pass, not the last.

func RegistrableDomain

func RegistrableDomain(host string) string

RegistrableDomain reduces a host to the domain a human would name as "the company's domain": the registrable label plus its public suffix. It is a pragmatic approximation, not a public-suffix implementation — it covers the multi-part suffixes that appear in practice and treats anything else as two-label. It never makes a security decision; it only groups hosts that probably belong to one operator.

A leading "www.", any port, and a trailing dot are removed, and the result is lowercased and punycoded. An IP literal is returned unchanged, because an address is never a company domain.

func StatusTextFor

func StatusTextFor(code int) string

StatusTextFor maps an HTTP status to a short, log-friendly explanation that distinguishes "nothing there" from "we were refused" from "they are busy".

func ValidatePublicURL

func ValidatePublicURL(u *url.URL) error

ValidatePublicURL reports whether a parsed URL is safe to fetch. It is the form callers use when they already hold a *url.URL, such as a link extracted from a page.

Types

type Budget

type Budget struct {
	// MaxPages caps documents fetched. Zero means unlimited.
	MaxPages int
	// MaxBytes caps total retained body bytes. Zero means unlimited.
	MaxBytes int64
	// MaxDuration caps wall-clock time. Zero means unlimited.
	MaxDuration time.Duration
	// MaxErrors caps consecutive failures before the crawl gives up on a host.
	MaxErrors int
	// MaxHostPages caps pages fetched from any single host, so one enormous
	// site cannot consume the whole run.
	MaxHostPages int
}

Budget is the cost envelope of a crawl. It is enforced, not advisory: the loop checks it before every fetch and stops when it is spent. This is the mechanism that keeps a run from quietly costing 400x its estimate.

type Config

type Config struct {
	// UserAgent identifies the crawler to servers. It must be honest: a real
	// contact address is part of being a good citizen and most operators
	// require one.
	UserAgent string

	// Concurrency is the number of in-flight fetches.
	Concurrency int
	// PerHostDelay is the minimum gap between two requests to the same host.
	// A site's Crawl-delay, when it declares one, overrides this upward.
	PerHostDelay time.Duration
	// FetchTimeout bounds a single request, including its body read.
	FetchTimeout time.Duration
	// MaxRetries is how many times a retryable failure is retried.
	MaxRetries int

	// MaxBytes caps a single response body. Larger responses are truncated and
	// flagged rather than downloaded in full.
	MaxBytes int64
	// MaxRedirects caps redirect chains. Zero means "do not follow".
	MaxRedirects int

	// RespectRobots disables nothing but is recorded in run metadata, because
	// turning it off is a decision an operator has to be able to audit.
	RespectRobots bool
	// RobotsCacheTTL is how long a host's robots.txt is trusted.
	RobotsCacheTTL time.Duration

	// ContentTypes lists the media types to retain.
	ContentTypes []string
	// DenyHostSuffixes blocks hosts by suffix, e.g. ".cdn.example" to skip an
	// asset host entirely.
	DenyHostSuffixes []string
	// AllowPrivateHosts permits fetching 10.x, 192.168.x, localhost and
	// link-local addresses. It is off by default and exists for local fixtures
	// and self-hosted targets; enabling it in production is an SSRF risk.
	AllowPrivateHosts bool

	// RawRetention is how long a stored body is kept. After that only
	// metadata, the hash and the extracted facts remain. Zero disables body
	// retention entirely, which is the right setting for a crawl whose output
	// is facts rather than archives.
	RawRetention time.Duration
	// KeepBodyInMemory controls whether Document.Body is populated. The store
	// decides separately whether to persist it.
	KeepBodyInMemory bool

	// MaxCrawlTime bounds a whole Crawl call.
	MaxCrawlTime time.Duration
	// MaxPages bounds a whole Crawl call.
	MaxPages int

	// LinkFilter decides whether a discovered URL is worth enqueuing. It runs
	// after the host policy. Returning false for a URL is a policy decision,
	// not an error.
	LinkFilter func(u *url.URL) bool
}

Config is the crawler's static configuration. It is validated once at construction; an invalid crawler is never built.

func DefaultConfig

func DefaultConfig() Config

DefaultConfig returns a configuration that is safe to run against the public internet: one connection, one request per second per host, robots honoured.

func ScreenConfig

func ScreenConfig() Config

ScreenConfig tightens the defaults for the cheap pre-filter pass: fewer bytes, no retention, no link following beyond a single hop.

func (Config) AllowsContentType

func (c Config) AllowsContentType(mediaType string) bool

AllowsContentType reports whether a media type is retained.

func (Config) AllowsHost

func (c Config) AllowsHost(u *url.URL) bool

AllowsHost reports whether a URL passes the suffix deny list. An empty host is rejected: a relative or malformed URL must never reach the network.

func (Config) Validate

func (c Config) Validate() error

Validate reports every problem at once.

type CrawlOptions

type CrawlOptions struct {
	// Seeds are the URLs to start from. They are normalised, checked against
	// host policy, and enqueued at the highest priority.
	Seeds []string
	// Sitemaps are sitemap URLs read for additional seeds.
	Sitemaps []string
	// RunID scopes stored rows. Empty generates one.
	RunID string
	// Depth is the research depth; it caps how far links are followed.
	Depth Depth
	// Budget caps the run. Zero fields fall back to the limits in Config.
	Budget Budget
	// FollowLinks enables link discovery. Off means the seeds are fetched and
	// nothing else, which is the cheap path a screening pass uses.
	FollowLinks bool
	// OnDocument is called for every stored document, inline. It must not
	// block; the pipeline passes an enqueue, not a network call.
	OnDocument func(Document) error
	// OnError is called for every non-fatal fetch failure.
	OnError func(FetchError)
	// PollInterval is how long Claim waits before re-checking when the frontier
	// has nothing claimable but workers are still running. Zero selects 250ms.
	PollInterval time.Duration
}

CrawlOptions configure one Crawl call.

type Crawler

type Crawler struct {
	// contains filtered or unexported fields
}

Crawler is the frontier-driven loop. It is safe for concurrent use.

func New

func New(fetch Fetcher, frontier Frontier, store Store, cfg Config, opt Options) (*Crawler, error)

New builds a Crawler over the given dependencies. Nothing is global: a caller wanting a second independent crawler passes a second Store and Frontier.

func (*Crawler) Crawl

func (c *Crawler) Crawl(ctx context.Context, opts CrawlOptions) (Result, error)

Crawl runs until the frontier drains or a budget is spent.

type Depth

type Depth int

Depth is a research depth. The same crawl is run at different depths depending on how much a candidate is worth: screen everything cheaply, then spend on the survivors.

const (
	// DepthScreen fetches only the URLs a source handed us. One page, no
	// link following, no sitemap expansion. This is the filter that runs
	// before any spending.
	DepthScreen Depth = iota
	// DepthStandard follows in-site links to a bounded depth and expands
	// sitemaps, which is enough to find a company page, its about page and its
	// contact page.
	DepthStandard
	// DepthDeep follows links more widely, expands sitemaps aggressively, and
	// fetches every discovered page kind that tends to carry facts (team,
	// careers, pricing, news, changelog).
	DepthDeep
)

func ParseDepth

func ParseDepth(s string) (Depth, error)

ParseDepth reads a depth name.

func (Depth) MaxDepth

func (d Depth) MaxDepth() int

MaxDepthFor is the hop limit a depth implies.

func (Depth) PageBudgetFor

func (d Depth) PageBudgetFor() int

PageBudgetFor is the default page ceiling a depth implies, before the caller's own budget is applied.

func (Depth) String

func (d Depth) String() string

String renders the depth for logs and run metadata.

type Document

type Document struct {
	// ID is assigned by the Store. Empty for an unsaved document.
	ID string
	// URL is the final URL after redirects — what the bytes actually are.
	URL string
	// RequestedURL is what we asked for. Kept so a redirect can be reported
	// and so a dedupe key can distinguish the request from the destination.
	RequestedURL string
	// RedirectChain records each hop, in order.
	RedirectChain []string

	// Status is the final HTTP status.
	Status int
	// StatusText explains a non-2xx outcome in words, for humans and logs.
	StatusText string
	// ContentType is the media type without parameters.
	ContentType string
	// Charset is the declared or detected character set, normalised to UTF-8
	// when we were able to convert.
	Charset string
	// ContentLength is the declared length, or -1 when unknown.
	ContentLength int64

	// ETag and LastModified drive conditional re-fetch, which is how a repeat
	// crawl of an unchanged site costs almost nothing.
	ETag         string
	LastModified time.Time

	// Body is the decoded response body, truncated at the configured cap.
	Body []byte
	// BodyTruncated reports that the cap was hit and the body is incomplete.
	BodyTruncated bool
	// ContentHash addresses the normalised body. Two documents with the same
	// hash are the same page as far as every downstream module is concerned.
	ContentHash string
	// RawHash is the digest of the bytes exactly as received, before charset
	// conversion. It is what conditional GET and dedupe-on-receipt use.
	RawHash string

	// Links are the links discovered in this document.
	Links []Link
	// Depth is the hop count from the seed.
	Depth int
	// FetchedAt is when the response was received.
	FetchedAt time.Time
	// NotModified reports a 304: the stored copy is still current.
	NotModified bool
	// Filtered reports that the response was received but its content type is
	// not one this crawler retains, so no body was kept.
	Filtered bool
	// FilterReason explains a Filtered document, and a skip, in one line. It is
	// empty for a normal fetch. Without it a caller cannot tell a deliberate
	// policy decision from a bug.
	FilterReason string
	// Elapsed is the wall-clock time the request took.
	Elapsed time.Duration
	// Attempt is the 1-based retry attempt that produced this document.
	Attempt int
	// UserAgent is the agent string used.
	UserAgent string
	// SourceHint records what discovered this URL, for attribution back to a
	// seed, a sitemap, or a link.
	SourceHint string
	// RunID scopes the document to a crawl, so a run's cost can be attributed.
	RunID string
}

Document is a fetched page. It is the crawler's product and the only thing other modules need to know about it.

Body is populated only when the caller asked for it or when the store is configured to retain raw bytes. Everything else is metadata that is cheap to keep forever.

type FetchError

type FetchError struct {
	URL     string
	Status  int
	Reason  string
	Attempt int
	Err     error
}

FetchError records a failed fetch without aborting a crawl.

func (FetchError) Error

func (e FetchError) Error() string

func (FetchError) Unwrap

func (e FetchError) Unwrap() error

type FetchResult

type FetchResult struct {
	Document Document
	Outcome  Outcome
	// Err is non-nil only for OutcomeFailed.
	Err error
	// Retryable reports whether the caller should try again later.
	Retryable bool
}

FetchResult is a Fetcher's answer: a document, a classification, and an error. The document is returned even on failure whenever there is something worth recording (a 404 status, a truncated body).

type Fetcher

type Fetcher interface {
	Fetch(ctx context.Context, req Request) FetchResult
}

Fetcher retrieves one URL. It holds no queue and no crawl state, so it can be used on its own as a "fetch this URL politely" library.

type Frontier

type Frontier interface {
	// Enqueue adds URLs, ignoring any whose dedupe key is already present.
	// It returns how many were newly added.
	Enqueue(ctx context.Context, items []Item) (int, error)
	// Claim leases up to n items, highest priority first.
	Claim(ctx context.Context, n int, lease time.Duration) ([]Item, error)
	// Complete marks items done. Items that fail are released with Backoff.
	Complete(ctx context.Context, items []Item, err error) error
	// Depth reports how many items are waiting.
	Depth(ctx context.Context) (int, error)
	// Pending reports how long until the next item becomes claimable, and
	// whether anything is pending at all. An exact answer lets a supervisor
	// sleep precisely as long as a backoff needs instead of polling, and lets it
	// stop the instant the queue is truly finished. A poll-count heuristic
	// cannot: an empty claim looks the same whether the crawl is done or a
	// worker still holds the only remaining item.
	Pending(ctx context.Context) (time.Duration, bool, error)
	// Reset clears a run's queue.
	Reset(ctx context.Context, runID string) error
	// Close releases resources.
	Close()
}

Frontier is the crawl queue. Claim/complete rather than pop/push, because a worker that dies mid-fetch must not lose its work: a claimed item that is never completed becomes claimable again once its lease expires.

type HTMLExtractor

type HTMLExtractor struct {
	// MaxLinks caps the returned set. Zero means DefaultMaxLinks.
	MaxLinks int
	// IncludeNoFollow keeps rel=nofollow links. They are links a site chose
	// not to endorse, but they are still real pages; the caller decides.
	IncludeNoFollow bool
}

HTMLExtractor extracts links from HTML, plus link targets worth fetching even when a page does not link to them in the body.

func (e HTMLExtractor) Links(doc Document, max int) []Link

Links implements LinkExtractor.

type HTTPFetcher

type HTTPFetcher struct {
	// contains filtered or unexported fields
}

HTTPFetcher is the default Fetcher.

func NewHTTPFetcher

func NewHTTPFetcher(cfg Config, opt Options) (*HTTPFetcher, error)

NewHTTPFetcher builds a Fetcher. It fails on an invalid config rather than silently crawling with nonsensical limits.

func (*HTTPFetcher) Fetch

func (f *HTTPFetcher) Fetch(ctx context.Context, req Request) FetchResult

Fetch retrieves one URL with robots compliance, politeness, conditional re-fetch, size and content-type policy, and bounded retries.

func (*HTTPFetcher) RobotsFor

func (f *HTTPFetcher) RobotsFor(origin string) (*robotstxt.Robot, bool)

RobotsFor returns the rules cached for a host, and whether any were cached. It is what a debug endpoint needs to answer "why was this URL skipped?".

type Item

type Item struct {
	// URL is the normalised absolute URL to fetch.
	URL string
	// Key is the dedupe identity: the URL without tracking parameters.
	Key string
	// Depth is the hop count from the seed.
	Depth int
	// Priority orders the queue. Higher is sooner.
	Priority int
	// SourceHint records what discovered the URL.
	SourceHint string
	// EnqueuedAt is when it entered the queue.
	EnqueuedAt time.Time
	// Attempts counts how many times it has been claimed.
	Attempts int
	// RunID scopes the item to a crawl.
	RunID string
}

Item is a unit of work in the frontier.

type Link struct {
	// Href is the absolute, normalised URL.
	Href string
	// Text is the anchor text, whitespace-collapsed and length-capped.
	Text string
	// Rel is the anchor's rel attribute, which distinguishes navigational
	// links from stylesheets and social profiles.
	Rel string
	// Internal reports whether the link stays on the document's own host.
	Internal bool
	// NoFollow reports rel="nofollow", a signal not to treat the target as
	// endorsed.
	NoFollow bool
}

Link is an outgoing hyperlink discovered in a document.

type LinkExtractor

type LinkExtractor interface {
	// Links returns the outgoing links of a document, absolute and
	// deduplicated, capped at max.
	Links(doc Document, max int) []Link
}

LinkExtractor pulls outgoing links out of a document. It is its own type so the HTML parsing rules can be unit tested without a network, and so a future extractor module can supply a richer implementation behind the same interface.

type Logger

type Logger interface {
	Debug(msg string, args ...any)
	Warn(msg string, args ...any)
}

Logger is the minimal logging surface the crawler needs, so the module does not force a particular logger on its consumers.

type MemoryFrontier

type MemoryFrontier struct {
	// contains filtered or unexported fields
}

MemoryFrontier is an in-process priority frontier. It backs unit tests, the one-shot CLI, and a single-worker crawl that does not need durability.

func NewMemoryFrontier

func NewMemoryFrontier(maxSize int) *MemoryFrontier

NewMemoryFrontier returns an empty frontier. maxSize bounds growth; zero selects DefaultFrontierSize.

func (*MemoryFrontier) Claim

func (f *MemoryFrontier) Claim(_ context.Context, n int, lease time.Duration) ([]Item, error)

Claim implements Frontier. Items whose lease has expired are eligible again.

func (*MemoryFrontier) Close

func (f *MemoryFrontier) Close()

Close implements Frontier.

func (*MemoryFrontier) Complete

func (f *MemoryFrontier) Complete(_ context.Context, items []Item, err error) error

Complete implements Frontier.

func (*MemoryFrontier) Depth

func (f *MemoryFrontier) Depth(context.Context) (int, error)

Depth implements Frontier. It counts every item still held, including leased ones, so a caller can tell "queue empty" from "queue leased out".

func (*MemoryFrontier) Enqueue

func (f *MemoryFrontier) Enqueue(_ context.Context, items []Item) (int, error)

Enqueue implements Frontier.

func (*MemoryFrontier) Evictions

func (f *MemoryFrontier) Evictions() int

Evictions reports how many queued items were dropped to make room. A non-zero value means the cap was hit and the crawl is narrower than intended.

func (*MemoryFrontier) Pending

func (f *MemoryFrontier) Pending(_ context.Context) (time.Duration, bool, error)

Pending implements Frontier.

func (*MemoryFrontier) Reset

func (f *MemoryFrontier) Reset(_ context.Context, runID string) error

Reset implements Frontier.

func (*MemoryFrontier) SetClock

func (f *MemoryFrontier) SetClock(now func() time.Time)

SetClock replaces the time source, for deterministic tests.

type MemoryStore

type MemoryStore struct {
	// contains filtered or unexported fields
}

MemoryStore keeps documents in process memory. It backs tests and short one-shot runs; nothing survives a restart.

It is safe for concurrent use: a Crawler's workers save in parallel, and an unguarded map write here is a real crash under load, not a theoretical one.

func NewMemoryStore

func NewMemoryStore() *MemoryStore

NewMemoryStore returns an empty in-memory store.

func (*MemoryStore) All

func (s *MemoryStore) All() []Document

All returns every stored document, in insertion order. Tests use it to assert what a crawl actually produced; it is not part of Store.

func (*MemoryStore) Close

func (s *MemoryStore) Close()

Close implements Store.

func (*MemoryStore) Get

func (s *MemoryStore) Get(_ context.Context, id string) (*Document, error)

Get implements Store.

func (*MemoryStore) LastDocument

func (s *MemoryStore) LastDocument(_ context.Context, url string) (*Document, error)

LastDocument implements Store.

func (*MemoryStore) Save

func (s *MemoryStore) Save(_ context.Context, doc Document) (Document, error)

Save implements Store.

func (*MemoryStore) Seen

func (s *MemoryStore) Seen(_ context.Context, contentHash string) (bool, error)

Seen implements Store.

func (*MemoryStore) SetClock

func (s *MemoryStore) SetClock(now func() time.Time)

SetClock replaces the time source, for deterministic tests.

func (*MemoryStore) Stats

func (s *MemoryStore) Stats(context.Context, string) (Stats, error)

Stats implements Store.

type Metrics

type Metrics struct {
	// contains filtered or unexported fields
}

Metrics wraps the platform registry with the crawler's own series. A nil *Metrics is valid and does nothing, so the crawler can run with no telemetry.

func NewMetrics

func NewMetrics(reg *observe.Registry) *Metrics

NewMetrics registers the crawler's series on reg.

type Options

type Options struct {
	// Store supplies the previous ETag/Last-Modified for a URL so a repeat
	// fetch can be conditional. Nil disables conditional requests.
	Store Store
	// UserAgent overrides Config.UserAgent for this fetcher.
	UserAgent string
	// HTTPClient overrides the constructed client. Tests use this; production
	// leaves it nil so the crawler builds a client with the right transport.
	HTTPClient *http.Client
	// Clock overrides time.Now, for deterministic tests.
	Clock func() time.Time
	// Limiter overrides the per-host rate limiter.
	Limiter *ratelimit.Keyed
	// Metrics receives counters. Nil disables metrics.
	Metrics *Metrics
	// Log is optional.
	Log Logger
}

Options configure a Fetcher.

type Outcome

type Outcome string

Outcome classifies a fetch so callers can branch without string matching.

const (
	// OutcomeFetched: a 2xx response with a usable body.
	OutcomeFetched Outcome = "fetched"
	// OutcomeNotModified: 304; the stored copy is still current.
	OutcomeNotModified Outcome = "not_modified"
	// OutcomeSkipped: policy declined to fetch (robots, type, scheme, budget).
	OutcomeSkipped Outcome = "skipped"
	// OutcomeFailed: the fetch was attempted and failed.
	OutcomeFailed Outcome = "failed"
	// OutcomeFiltered: the document was fetched but its content is not
	// something this crawler stores or links from.
	OutcomeFiltered Outcome = "filtered"
)

type PostgresFrontier

type PostgresFrontier struct {
	// contains filtered or unexported fields
}

PostgresFrontier is a durable, shareable Frontier. Several crawler processes can claim from the same queue, and a process that dies mid-fetch leaves a lease that expires rather than a permanently stuck item.

func NewPostgresFrontier

func NewPostgresFrontier(pool *pgxpool.Pool) *PostgresFrontier

NewPostgresFrontier wraps a pool.

func (*PostgresFrontier) Claim

func (f *PostgresFrontier) Claim(ctx context.Context, n int, lease time.Duration) ([]Item, error)

Claim implements Frontier with FOR UPDATE SKIP LOCKED, so two workers never lease the same item. An item whose lease has expired is claimable again.

func (*PostgresFrontier) Close

func (f *PostgresFrontier) Close()

Close implements Frontier. It is a no-op: the pool belongs to the process, and several frontiers can share it.

func (*PostgresFrontier) Complete

func (f *PostgresFrontier) Complete(ctx context.Context, items []Item, err error) error

Complete implements Frontier. A nil error deletes the item; an error backs it off exponentially, and an item that keeps failing is dropped so a crawl is not a queue with infinite patience.

func (*PostgresFrontier) Depth

func (f *PostgresFrontier) Depth(ctx context.Context) (int, error)

Depth implements Frontier.

func (*PostgresFrontier) Enqueue

func (f *PostgresFrontier) Enqueue(ctx context.Context, items []Item) (int, error)

Enqueue implements Frontier.

The dedupe gate is a single statement: a CTE inserts the run's key into crawler_frontier_seen and queues the URL only if that insert was new. Doing it as two statements would require deleting the queue row on a conflict, and that delete would also remove a legitimately queued item for the same key.

func (*PostgresFrontier) Pending

func (f *PostgresFrontier) Pending(ctx context.Context) (time.Duration, bool, error)

Pending implements Frontier. It is one query rather than a count plus a scan, so a supervisor does not have to guess how long to sleep.

func (*PostgresFrontier) Reset

func (f *PostgresFrontier) Reset(ctx context.Context, runID string) error

Reset implements Frontier. With an empty runID it clears everything, which is what a fresh crawl wants; with a runID it clears only that run's queue and its seen keys, so a concurrent run is untouched.

func (*PostgresFrontier) SetClock

func (f *PostgresFrontier) SetClock(now func() time.Time)

SetClock replaces the time source, for deterministic tests.

type PostgresStore

type PostgresStore struct {
	// contains filtered or unexported fields
}

PostgresStore is the durable Store. It is what lets a crawl be resumable and what makes conditional re-fetch possible across restarts.

func NewPostgresStore

func NewPostgresStore(pool *pgxpool.Pool) *PostgresStore

NewPostgresStore wraps a pool. The pool is not owned by the store: the process that built it closes it.

func (*PostgresStore) Close

func (s *PostgresStore) Close()

Close implements Store. It closes the pool, because the store is the only holder of it in a standalone crawler process.

func (*PostgresStore) Get

func (s *PostgresStore) Get(ctx context.Context, id string) (*Document, error)

Get implements Store.

func (*PostgresStore) LastDocument

func (s *PostgresStore) LastDocument(ctx context.Context, url string) (*Document, error)

LastDocument implements Store.

func (*PostgresStore) Pool

func (s *PostgresStore) Pool() *pgxpool.Pool

Pool exposes the underlying pool for callers that need to run their own queries against crawler tables.

func (*PostgresStore) RecordRun

func (s *PostgresStore) RecordRun(ctx context.Context, r Result) error

RecordRun writes a completed run's accounting.

func (*PostgresStore) Save

func (s *PostgresStore) Save(ctx context.Context, doc Document) (Document, error)

Save implements Store. A repeat fetch of the same URL updates the existing row in place, so a daily crawl of a static site does not grow the table.

func (*PostgresStore) Seen

func (s *PostgresStore) Seen(ctx context.Context, contentHash string) (bool, error)

Seen implements Store.

func (*PostgresStore) Stats

func (s *PostgresStore) Stats(ctx context.Context, runID string) (Stats, error)

Stats implements Store.

type Request

type Request struct {
	// URL is the absolute URL to fetch.
	URL string
	// Depth is the hop count from the seed, used for budget decisions.
	Depth int
	// Priority orders the frontier. Higher is fetched sooner.
	Priority int
	// ETag and IfModifiedSince enable a conditional request when set.
	ETag            string
	IfModifiedSince time.Time
	// SourceHint records what discovered this URL, for attribution.
	SourceHint string
}

Request asks the Fetcher for one URL.

type Result

type Result struct {
	RunID         string        `json:"run_id"`
	Seeds         int           `json:"seeds"`
	Enqueued      int           `json:"enqueued"`
	Fetched       int           `json:"fetched"`
	NotModified   int           `json:"not_modified"`
	Skipped       int           `json:"skipped"`
	Filtered      int           `json:"filtered"`
	Failed        int           `json:"failed"`
	LinksSeen     int           `json:"links_seen"`
	BytesRetained int64         `json:"bytes_retained"`
	Stats         Stats         `json:"stats"`
	Duration      time.Duration `json:"-"`
	Errors        []FetchError  `json:"-"`
	// contains filtered or unexported fields
}

Result summarises a finished crawl. It is valid even when Crawl returns ErrCrawlStopped.

type RobotsReporter

type RobotsReporter interface {
	RobotsFor(origin string) (*robotstxt.Robot, bool)
}

RobotsReporter is the optional capability of a Fetcher that can report the rules it holds for a host. It is separate from Fetcher because a caller with its own transport has no reason to implement it, and forcing the method on the interface would make that awkward.

type RunState

type RunState struct {
	StartedAt     time.Time
	PagesFetched  int
	BytesRetained int64
	PerHost       map[string]int
	Consecutive   map[string]int
}

RunState is the mutable accounting for one crawl. It is owned by the Crawler and guarded by its mutex; Budget reads it under that lock.

type SitemapEntry

type SitemapEntry struct {
	Loc        string
	LastMod    string
	ChangeFreq string
	Priority   string
}

SitemapEntry is one <url> in a sitemap.

type SitemapIndexEntry

type SitemapIndexEntry struct {
	Loc     string
	LastMod string
}

SitemapIndexEntry is one <sitemap> in a sitemap index.

type SlogLogger

type SlogLogger struct{ L *slog.Logger }

SlogLogger adapts a *slog.Logger to the crawler's Logger interface.

func (SlogLogger) Debug

func (s SlogLogger) Debug(msg string, args ...any)

func (SlogLogger) Warn

func (s SlogLogger) Warn(msg string, args ...any)

type Stats

type Stats struct {
	RunID          string    `json:"run_id,omitempty"`
	Documents      int       `json:"documents"`
	BytesRetained  int64     `json:"bytes_retained"`
	Skipped        int       `json:"skipped"`
	Failed         int       `json:"failed"`
	NotModified    int       `json:"not_modified"`
	Filtered       int       `json:"filtered"`
	StartedAt      time.Time `json:"started_at,omitempty"`
	FinishedAt     time.Time `json:"finished_at,omitempty"`
	DurationMillis int64     `json:"duration_ms"`
	StoppedBy      string    `json:"stopped_by,omitempty"`
}

Stats summarises a crawl.

type Store

type Store interface {
	// Save writes a document and returns it with ID assigned. A document whose
	// content hash already exists for that URL updates the existing row rather
	// than inserting a duplicate, but still records the new fetch.
	Save(ctx context.Context, doc Document) (Document, error)
	// LastDocument returns the most recent stored document for a URL, or nil
	// when the URL has never been fetched.
	LastDocument(ctx context.Context, url string) (*Document, error)
	// Get returns a stored document by ID.
	Get(ctx context.Context, id string) (*Document, error)
	// Seen reports whether a content hash is already known, which is how a
	// repeated crawl of a static site costs one conditional request per URL
	// rather than one row per URL per day.
	Seen(ctx context.Context, contentHash string) (bool, error)
	// Stats summarises a crawl for the API and for logs.
	Stats(ctx context.Context, runID string) (Stats, error)
	// Close releases resources.
	Close()
}

Store persists fetched documents. It is the crawler's only durable state and the only thing the Fetcher reads to make a conditional request.

Directories

Path Synopsis
cmd
crawld command
Command crawld serves the crawler as a standalone service.
Command crawld serves the crawler as a standalone service.
internal
Package robotstxt implements the Robots Exclusion Protocol as specified in RFC 9309, plus the long-standing crawl-delay and sitemap extensions.
Package robotstxt implements the Robots Exclusion Protocol as specified in RFC 9309, plus the long-standing crawl-delay and sitemap extensions.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL