// Package lookup looks up each client's AS number and country, through // the GeoJS web service or in the lookup database, the IPinfo Lite file // SWWAF_LOOKUP_DB_PATH names. GeoJS's answers are kept in memory, for at // most 100,000 clients and for 7 days each, and are written to // lookups.json and read from it by the state package. package lookup import ( "context" "encoding/json" "errors" "fmt" "io" "log/slog" "net/http" "net/netip" "slices" "strconv" "strings" "sync" "time" "github.com/hashicorp/golang-lru/v2/simplelru" "sneak.berlin/go/smallwebwaf/internal/alerts" "sneak.berlin/go/smallwebwaf/internal/metrics" ) // URL is GeoJS's endpoint for an address's place and network. Asked about // several addresses at once, comma separated in its ip parameter, it // answers with a list. const URL = "https://get.geojs.io/v1/ip/geo.json" const ( // keepFor is how long an answer is used instead of asking GeoJS again. keepFor = 7 * 24 * time.Hour // maxAnswers is how many answers are kept. Past it, the one used // longest ago is dropped. maxAnswers = 100000 // maxWaiting is how many clients may wait to be asked about. Past it, // a new client counts as not found and is not asked about until there // is room, so that a swarm of new addresses while GeoJS is down cannot // fill the memory. maxWaiting = 10000 // maxPerRequest is how many addresses one request to GeoJS asks about. maxPerRequest = 200 // unknownASN is the AS number GeoJS gives when it knows none. unknownASN = 64512 // After a failure GeoJS is not asked again for a second, and for // retryDelayFactor times as long after each further failure in a row, // up to five minutes. firstRetryDelay = time.Second retryDelayFactor = 2 maxRetryDelay = 5 * time.Minute // maxResponseBytes is the most of GeoJS's answer that is read. maxResponseBytes = 1 << 20 ) var ( errStatus = errors.New("GeoJS answered") errLeftOut = errors.New("GeoJS's answer left out") ) // Params are what New needs. type Params struct { // URL is where GeoJS is asked, normally URL. URL string // Timeout is how long a request waits for its client's first answer, // and how long a request to GeoJS may take before it is abandoned // (SWWAF_LOOKUP_TIMEOUT). Timeout time.Duration // Wait is true when a setting needs each request's answer before the // request goes on. Otherwise no request waits for one. Wait bool // Answered, unless nil, is given each answer GeoJS gives, once it is // kept. Answered func(Answer) // Now tells the time, normally time.Now. Now func() time.Time // ProcessLog receives GeoJS's failures. ProcessLog *slog.Logger // Metrics count the requests to GeoJS, those that failed, and the // clients that go without an answer. Metrics *metrics.Metrics // Alerts receive a source_failure alert each time GeoJS fails. Alerts *alerts.Queue } // GeoJS looks up clients' AS numbers and countries through GeoJS. At most // one request to GeoJS is under way at a time, and it asks about every // client waiting, up to maxPerRequest. It is safe for concurrent use. type GeoJS struct { url string timeout time.Duration wait bool answered func(Answer) now func() time.Time processLog *slog.Logger metrics *metrics.Metrics alerts *alerts.Queue // httpClient follows no redirect, so that visitors' addresses go to // GeoJS alone: a redirect is a failure. httpClient *http.Client mu sync.Mutex answers *simplelru.LRU[netip.Prefix, *Answer] // waiting are the clients without an answer: those to ask GeoJS about, // and those it is being asked about. waiting map[netip.Prefix]*wait // asking is true while a request to GeoJS is under way. asking bool // retryDelay is how long GeoJS is left alone after its last failure, // zero after an answer; retryAt is when it may be asked again. retryDelay time.Duration retryAt time.Time } // Answer is what GeoJS or the lookup database said about a client: its AS // number, such as AS64496, and the AS's name, both "" when the source knows // no AS number for it; its country, "" when the source cannot place it; // when the source said so; and, for GeoJS's answers, which lookups.json // holds, when the answer was last used. The zero Answer is that of a // client with no answer. // //nolint:tagliatelle // the state files use snake_case, as the request log does type Answer struct { Client netip.Prefix `json:"client"` ASN string `json:"asn"` ASName string `json:"as_name"` Country string `json:"country"` Answered time.Time `json:"answered"` Used time.Time `json:"used"` } // wait is a client waiting for its answer. type wait struct { // asked is closed when the client gets its answer, and closed and // replaced each time GeoJS fails before then. asked chan struct{} // late is true once the client has gone without an answer, for a // whole timeout or because GeoJS failed: its requests no longer wait. late bool } // New returns a GeoJS with no answer kept yet. func New(params Params) *GeoJS { answers, err := simplelru.NewLRU[netip.Prefix, *Answer](maxAnswers, nil) if err != nil { panic(err) // NewLRU fails only for a size below one } return &GeoJS{ url: params.URL, timeout: params.Timeout, wait: params.Wait, answered: params.Answered, now: params.Now, processLog: params.ProcessLog, metrics: params.Metrics, alerts: params.Alerts, httpClient: &http.Client{ CheckRedirect: func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse }, }, answers: answers, waiting: map[netip.Prefix]*wait{}, } } // LookUp returns the answer GeoJS gave about client, with its country as // a two-letter code in capitals, or the zero Answer when there is none // yet. An answer is kept for 7 days. Without one, the client is asked // about in the background, and, while Wait is set, the request waits up // to Timeout for the answer, unless the client has gone without one // before. ctx is the context of the client's request, and ends the wait // when it ends. // // GeoJS is asked about the client's first address, which is the client's // own address for IPv4, and an address in the same place for an IPv6 /64. func (g *GeoJS) LookUp(ctx context.Context, client netip.Prefix) Answer { answer, asked := g.answerOrWait(ctx, client) if asked == nil { return answer } timer := time.NewTimer(g.timeout) defer timer.Stop() select { case <-asked: case <-timer.C: case <-ctx.Done(): } g.mu.Lock() defer g.mu.Unlock() answer, found := g.kept(client) if !found { g.metrics.GeoJSUnanswered.Inc() } w, waiting := g.waiting[client] if !found && waiting { w.late = true } return answer } // Kept returns client's answer, if one is kept, without asking GeoJS. func (g *GeoJS) Kept(client netip.Prefix) (Answer, bool) { g.mu.Lock() defer g.mu.Unlock() return g.kept(client) } // Snapshot returns every answer kept, sorted by client, as lookups.json // lists them. func (g *GeoJS) Snapshot() []Answer { g.mu.Lock() answers := make([]Answer, 0, g.answers.Len()) for _, kept := range g.answers.Values() { answers = append(answers, *kept) } g.mu.Unlock() slices.SortFunc(answers, func(a, b Answer) int { return a.Client.Compare(b.Client) }) return answers } // Load keeps answers read from lookups.json, in place of the answers it // keeps, in the order they were last used, so that the one used longest // ago is dropped first. Answers GeoJS gave keepFor ago or more are // dropped. func (g *GeoJS) Load(answers []Answer) { answers = slices.Clone(answers) slices.SortStableFunc(answers, func(a, b Answer) int { return a.Used.Compare(b.Used) }) g.mu.Lock() defer g.mu.Unlock() g.answers.Purge() now := g.now() for _, answer := range answers { if now.Sub(answer.Answered) < keepFor { g.answers.Add(answer.Client, &answer) } } } // answerOrWait returns client's kept answer if it has one. Otherwise it // puts the client among those waiting if there is room, has GeoJS asked // about them if it can be, and returns what to wait on for the answer, or // nil when there is nothing to wait for. func (g *GeoJS) answerOrWait( ctx context.Context, client netip.Prefix, ) (Answer, <-chan struct{}) { g.mu.Lock() defer g.mu.Unlock() answer, found := g.kept(client) if found { return answer, nil } w, waiting := g.waiting[client] if !waiting && len(g.waiting) < maxWaiting { w = &wait{asked: make(chan struct{})} g.waiting[client] = w } g.ask(ctx) if !g.wait { return Answer{}, nil // the answer is not needed before the request goes on } if w == nil { g.metrics.GeoJSUnanswered.Inc() return Answer{}, nil // too many clients wait already } if !g.asking { // GeoJS is left alone after a failure, so no answer can come. w.late = true } if w.late { g.metrics.GeoJSUnanswered.Inc() return Answer{}, nil } return Answer{}, w.asked } // kept returns client's answer, if GeoJS gave it less than keepFor ago, // and notes that it was used. func (g *GeoJS) kept(client netip.Prefix) (Answer, bool) { now := g.now() kept, found := g.answers.Get(client) if !found || now.Sub(kept.Answered) >= keepFor { return Answer{}, false } kept.Used = now return *kept, true } // ask starts asking GeoJS about the waiting clients, unless a request to // it is under way or it is left alone after a failure. The requests to // GeoJS are for every client waiting, so they go on when the client's // request whose ctx is given ends. func (g *GeoJS) ask(ctx context.Context) { if g.asking || g.now().Before(g.retryAt) { return } g.asking = true go g.askAboutWaiting(context.WithoutCancel(ctx)) } // askAboutWaiting asks GeoJS about the waiting clients, one request at a // time, until none is left or GeoJS fails. Each answer kept is given to // Answered, outside the lock, since Answered takes locks of its own. func (g *GeoJS) askAboutWaiting(ctx context.Context) { for { clients := g.nextClients() if len(clients) == 0 { return } given, err := g.request(ctx, clients) kept, answered := g.keep(clients, given, err) if g.answered != nil { for _, answer := range kept { g.answered(answer) } } if !answered { return } } } // nextClients returns up to maxPerRequest of the waiting clients. When // none is waiting, it returns none and notes that no request to GeoJS is // under way. func (g *GeoJS) nextClients() []netip.Prefix { g.mu.Lock() defer g.mu.Unlock() if len(g.waiting) == 0 { g.asking = false return nil } clients := make([]netip.Prefix, 0, min(len(g.waiting), maxPerRequest)) for client := range g.waiting { if len(clients) == maxPerRequest { break } clients = append(clients, client) } return clients } // keep notes how a request to GeoJS about clients ended, given being the // answer for each address GeoJS's answer names. It returns the answers it // kept, and reports whether GeoJS answered about all of the clients. Each // client whose address GeoJS's answer names gets its answer. An answer // that leaves an address out is a failure. After a failure GeoJS is left // alone for a while, and every client still waiting stops waiting and is // asked about once GeoJS is asked again. func (g *GeoJS) keep( clients []netip.Prefix, given map[netip.Addr]Answer, err error, ) ([]Answer, bool) { g.mu.Lock() defer g.mu.Unlock() now := g.now() kept := make([]Answer, 0, len(clients)) leftOut := 0 for _, client := range clients { answer, named := given[client.Addr()] if !named { leftOut++ continue } answer.Client, answer.Answered, answer.Used = client, now, now g.answers.Add(client, &answer) kept = append(kept, answer) close(g.waiting[client].asked) delete(g.waiting, client) } if err == nil && leftOut > 0 { err = fmt.Errorf("%w %d of %d addresses", errLeftOut, leftOut, len(clients)) } if err != nil { g.metrics.GeoJSFailures.Inc() g.retryDelay = min(max(retryDelayFactor*g.retryDelay, firstRetryDelay), maxRetryDelay) g.retryAt = now.Add(g.retryDelay) g.asking = false for _, w := range g.waiting { close(w.asked) w.asked = make(chan struct{}) w.late = true } g.processLog.Warn("asking GeoJS failed", "error", err.Error(), "asking_again_in", g.retryDelay.String()) g.alerts.Raise(alerts.Alert{ Event: alerts.EventSourceFailure, Reason: "asking GeoJS failed", Detail: map[string]any{ "source": "geojs", "error": err.Error(), "asking_again_in": g.retryDelay.String(), }, }) return kept, false } g.retryDelay = 0 return kept, true } // request asks GeoJS about clients in one request, and returns the answer // for each address GeoJS's answer names: its AS number and the AS's name, // both "" for the AS number 64512, which GeoJS gives when it knows none, // and its country, in capitals. func (g *GeoJS) request( ctx context.Context, clients []netip.Prefix, ) (map[netip.Addr]Answer, error) { addrs := make([]string, 0, len(clients)) for _, client := range clients { addrs = append(addrs, client.Addr().String()) } ctx, cancel := context.WithTimeout(ctx, g.timeout) defer cancel() req, err := http.NewRequestWithContext(ctx, http.MethodGet, g.url, http.NoBody) if err != nil { return nil, fmt.Errorf("make the request to GeoJS: %w", err) } req.URL.RawQuery = "ip=" + strings.Join(addrs, ",") g.metrics.GeoJSRequests.Inc() res, err := g.httpClient.Do(req) if err != nil { // Do's error names the URL, and so the visitors' addresses, which // are not to be logged: only what went wrong is kept. return nil, fmt.Errorf("ask GeoJS: %w", errors.Unwrap(err)) } defer func() { _ = res.Body.Close() }() if res.StatusCode != http.StatusOK { return nil, fmt.Errorf("%w %s", errStatus, res.Status) } //nolint:tagliatelle // GeoJS's own names var answers []struct { IP string `json:"ip"` ASN int64 `json:"asn"` ASName string `json:"organization_name"` CountryCode string `json:"country_code"` } err = json.NewDecoder(io.LimitReader(res.Body, maxResponseBytes)).Decode(&answers) if err != nil { return nil, fmt.Errorf("read GeoJS's answer: %w", err) } given := make(map[netip.Addr]Answer, len(answers)) for _, item := range answers { addr, err := netip.ParseAddr(item.IP) if err != nil { continue } answer := Answer{Country: strings.ToUpper(item.CountryCode)} if item.ASN != 0 && item.ASN != unknownASN { answer.ASN = "AS" + strconv.FormatInt(item.ASN, 10) answer.ASName = item.ASName } given[addr] = answer } return given, nil }