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
- Variables
- func CharsetOf(header string) string
- func ContentTypeOf(header string) string
- func DefaultContentTypes() []string
- func IsHTMLContentType(ct string) bool
- func IsRetryableStatus(code int) bool
- func IsSameSite(a, b *url.URL) bool
- func IsXMLContentType(ct string) bool
- func MigrationSource() db.Source
- func NewHTTPClient(cfg Config) *http.Client
- func NormalizeURL(u *url.URL) (string, error)
- func ParseSitemap(r io.Reader) (entries []SitemapEntry, index []SitemapIndexEntry, err error)
- func PublicURL(raw string) (*url.URL, error)
- func RegistrableDomain(host string) string
- func StatusTextFor(code int) string
- func ValidatePublicURL(u *url.URL) error
- type Budget
- type Config
- type CrawlOptions
- type Crawler
- type Depth
- type Document
- type FetchError
- type FetchResult
- type Fetcher
- type Frontier
- type HTMLExtractor
- type HTTPFetcher
- type Item
- type Link
- type LinkExtractor
- type Logger
- type MemoryFrontier
- func (f *MemoryFrontier) Claim(_ context.Context, n int, lease time.Duration) ([]Item, error)
- func (f *MemoryFrontier) Close()
- func (f *MemoryFrontier) Complete(_ context.Context, items []Item, err error) error
- func (f *MemoryFrontier) Depth(context.Context) (int, error)
- func (f *MemoryFrontier) Enqueue(_ context.Context, items []Item) (int, error)
- func (f *MemoryFrontier) Evictions() int
- func (f *MemoryFrontier) Pending(_ context.Context) (time.Duration, bool, error)
- func (f *MemoryFrontier) Reset(_ context.Context, runID string) error
- func (f *MemoryFrontier) SetClock(now func() time.Time)
- type MemoryStore
- func (s *MemoryStore) All() []Document
- func (s *MemoryStore) Close()
- func (s *MemoryStore) Get(_ context.Context, id string) (*Document, error)
- func (s *MemoryStore) LastDocument(_ context.Context, url string) (*Document, error)
- func (s *MemoryStore) Save(_ context.Context, doc Document) (Document, error)
- func (s *MemoryStore) Seen(_ context.Context, contentHash string) (bool, error)
- func (s *MemoryStore) SetClock(now func() time.Time)
- func (s *MemoryStore) Stats(context.Context, string) (Stats, error)
- type Metrics
- type Options
- type Outcome
- type PostgresFrontier
- func (f *PostgresFrontier) Claim(ctx context.Context, n int, lease time.Duration) ([]Item, error)
- func (f *PostgresFrontier) Close()
- func (f *PostgresFrontier) Complete(ctx context.Context, items []Item, err error) error
- func (f *PostgresFrontier) Depth(ctx context.Context) (int, error)
- func (f *PostgresFrontier) Enqueue(ctx context.Context, items []Item) (int, error)
- func (f *PostgresFrontier) Pending(ctx context.Context) (time.Duration, bool, error)
- func (f *PostgresFrontier) Reset(ctx context.Context, runID string) error
- func (f *PostgresFrontier) SetClock(now func() time.Time)
- type PostgresStore
- func (s *PostgresStore) Close()
- func (s *PostgresStore) Get(ctx context.Context, id string) (*Document, error)
- func (s *PostgresStore) LastDocument(ctx context.Context, url string) (*Document, error)
- func (s *PostgresStore) Pool() *pgxpool.Pool
- func (s *PostgresStore) RecordRun(ctx context.Context, r Result) error
- func (s *PostgresStore) Save(ctx context.Context, doc Document) (Document, error)
- func (s *PostgresStore) Seen(ctx context.Context, contentHash string) (bool, error)
- func (s *PostgresStore) Stats(ctx context.Context, runID string) (Stats, error)
- type Request
- type Result
- type RobotsReporter
- type RunState
- type SitemapEntry
- type SitemapIndexEntry
- type SlogLogger
- type Stats
- type Store
Constants ¶
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.
const DefaultMaxLinks = 500
DefaultMaxLinks bounds link extraction so a navigation-heavy page cannot dominate a crawl.
Variables ¶
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.
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.
var ErrFrontierClosed = errors.New("crawler: frontier is closed")
ErrFrontierClosed is returned once the frontier has been closed.
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 ContentTypeOf ¶
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 ¶
IsHTMLContentType reports whether a media type is worth parsing for links and facts.
func IsRetryableStatus ¶
IsRetryableStatus reports whether a status justifies an immediate retry. 4xx other than 429 will not change by asking again.
func IsSameSite ¶
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 ¶
IsXMLContentType reports whether a media type is a sitemap or feed.
func MigrationSource ¶
MigrationSource returns the crawler's migrations for db.Migrate.
func NewHTTPClient ¶
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 ¶
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 ¶
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 ¶
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 ¶
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 ¶
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 ¶
AllowsContentType reports whether a media type is retained.
func (Config) AllowsHost ¶
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.
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.
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 (Depth) PageBudgetFor ¶
PageBudgetFor is the default page ceiling a depth implies, before the caller's own budget is applied.
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 ¶
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.
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.
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 ¶
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 ¶
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 ¶
Claim implements Frontier. Items whose lease has expired are eligible again.
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) 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) 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) LastDocument ¶
LastDocument implements Store.
func (*MemoryStore) SetClock ¶
func (s *MemoryStore) SetClock(now func() time.Time)
SetClock replaces the time source, for deterministic tests.
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 ¶
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 ¶
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 ¶
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 ¶
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 ¶
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) LastDocument ¶
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 ¶
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.
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 ¶
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 ¶
SitemapEntry is one <url> in a sitemap.
type SitemapIndexEntry ¶
SitemapIndexEntry is one <sitemap> in a sitemap index.
type SlogLogger ¶
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.
Source Files
¶
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. |