Verified end-to-end on the local OSRM stack: - internal/traffic: 511NY OData client (lenient decode, raw archive), Overpass way matcher (class-window selection: highest-class road within 150m beats a closer service road), incident policy map, CSV builder - cmd/traffic: fetch | match | apply; apply copies the dataset (CoW) and re-customizes the traffic layer in ~15s (v26 customize has no --output-prefix; it writes contract files in place) - Proof: synthetic I-90 closure → 38s (static :5000) vs 142s (traffic :5002) on the same 0.87 km stretch - Policy speeds are 'assumed' provenance; R4: traffic routes must carry computed(traffic, as_of, policy=assumed)
237 lines
7.3 KiB
Go
237 lines
7.3 KiB
Go
// Command traffic runs the live-traffic pipeline:
|
|
//
|
|
// traffic fetch --key K [--out events.json] # pull 511NY incidents
|
|
// traffic match --events events.json [--out speeds.csv] # Overpass-matched node-pair speeds
|
|
// traffic apply --csv speeds.csv --data <dir> # osrm-customize with the speed file
|
|
//
|
|
// `fetch` needs a 511NY developer key (511ny.org → Sign Up → API key).
|
|
// Until a key is available, `match`/`apply` work from any events JSON in the
|
|
// shape FetchEvents emits (or by hand).
|
|
package main
|
|
|
|
import (
|
|
"encoding/json"
|
|
"flag"
|
|
"fmt"
|
|
"os"
|
|
"os/exec"
|
|
"path/filepath"
|
|
"strings"
|
|
"time"
|
|
|
|
"maps/router/internal/traffic"
|
|
)
|
|
|
|
func main() {
|
|
if len(os.Args) < 2 {
|
|
usage()
|
|
}
|
|
var err error
|
|
switch os.Args[1] {
|
|
case "fetch":
|
|
err = cmdFetch(os.Args[2:])
|
|
case "match":
|
|
err = cmdMatch(os.Args[2:])
|
|
case "apply":
|
|
err = cmdApply(os.Args[2:])
|
|
case "help", "-h", "--help":
|
|
usage()
|
|
default:
|
|
usage()
|
|
}
|
|
if err != nil {
|
|
fmt.Fprintln(os.Stderr, "error:", err)
|
|
os.Exit(1)
|
|
}
|
|
}
|
|
|
|
func usage() {
|
|
fmt.Fprint(os.Stderr, `traffic: live-traffic pipeline (511NY → OSM match → OSRM segment speeds)
|
|
|
|
traffic fetch --key K [--out events.json] [--raw out-raw.json]
|
|
traffic match --events events.json [--out speeds.csv] [--radius 3000]
|
|
traffic apply --csv speeds.csv --data /path/to/data
|
|
`)
|
|
os.Exit(2)
|
|
}
|
|
|
|
func cmdFetch(args []string) error {
|
|
fs := flag.NewFlagSet("fetch", flag.ExitOnError)
|
|
key := fs.String("key", envOr("FIFTYONE_NY_KEY", ""), "511NY developer API key")
|
|
out := fs.String("out", "events.json", "normalized events output")
|
|
raw := fs.String("raw", "", "archive raw response (optional)")
|
|
fs.Parse(args)
|
|
if *key == "" {
|
|
return fmt.Errorf("no --key (get one: 511ny.org → Sign Up → account → API key)")
|
|
}
|
|
events, body, err := traffic.FetchEvents(*key)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if *raw != "" {
|
|
if err := os.WriteFile(*raw, body, 0o644); err != nil {
|
|
return err
|
|
}
|
|
fmt.Fprintf(os.Stderr, "raw response → %s\n", *raw)
|
|
}
|
|
valid, matched := 0, 0
|
|
for _, e := range events {
|
|
if e.Valid() {
|
|
valid++
|
|
}
|
|
if p, ok := traffic.PolicyFor(e.Type); ok {
|
|
matched++
|
|
_ = p
|
|
}
|
|
}
|
|
data, _ := json.MarshalIndent(events, "", " ")
|
|
if err := os.WriteFile(*out, data, 0o644); err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("fetched %d events (%d located, %d with traffic policy) → %s\n", len(events), valid, matched, *out)
|
|
return nil
|
|
}
|
|
|
|
func cmdMatch(args []string) error {
|
|
fs := flag.NewFlagSet("match", flag.ExitOnError)
|
|
eventsPath := fs.String("events", "events.json", "normalized events JSON")
|
|
out := fs.String("out", "speeds.csv", "segment-speed CSV output")
|
|
radius := fs.Float64("radius", 3000, "Overpass search radius (m)")
|
|
overpass := fs.String("overpass", "", "Overpass endpoint (default: public)")
|
|
timeout := fs.Int("timeout", 15, "Overpass query timeout (s)")
|
|
fs.Parse(args)
|
|
|
|
raw, err := os.ReadFile(*eventsPath)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
var events []traffic.Event
|
|
if err := json.Unmarshal(raw, &events); err != nil {
|
|
return fmt.Errorf("events json: %w", err)
|
|
}
|
|
|
|
var assignments []traffic.Assignment
|
|
skipped, noPolicy := 0, 0
|
|
for i, e := range events {
|
|
if !e.Valid() {
|
|
skipped++
|
|
continue
|
|
}
|
|
p, ok := traffic.PolicyFor(e.Type)
|
|
if !ok {
|
|
noPolicy++
|
|
continue
|
|
}
|
|
ways, err := traffic.QueryWaysAround(nil, *overpass, e.Lat, e.Lon, *radius, *timeout)
|
|
if err != nil {
|
|
fmt.Fprintf(os.Stderr, "event %d (%s): %v\n", i, e.Type, err)
|
|
skipped++
|
|
continue
|
|
}
|
|
way, a, b, dist, ok := traffic.ClosestNodePair(ways, e.Lat, e.Lon)
|
|
if !ok {
|
|
fmt.Fprintf(os.Stderr, "event %d (%s): no highway way found\n", i, e.Type)
|
|
skipped++
|
|
continue
|
|
}
|
|
pairs := traffic.ExpandToRadius(way, a, b, p.RadiusM)
|
|
if len(pairs) == 0 {
|
|
pairs = [][2]int64{{a, b}}
|
|
}
|
|
fmt.Fprintf(os.Stderr, "event %d: %q @ %s (%s, %dm to %s/%s) → %d edges @ %d km/h\n",
|
|
i, e.Type, e.City, way.Ref, int(dist), way.Highway, e.Type, len(pairs), p.SpeedKmh)
|
|
for _, pr := range pairs {
|
|
assignments = append(assignments, traffic.Assignment{Pair: pr, SpeedKmh: p.SpeedKmh, Weight: p.Weight})
|
|
}
|
|
// be gentle with the public Overpass instance
|
|
time.Sleep(500 * time.Millisecond)
|
|
}
|
|
|
|
csv := traffic.BuildCSV(assignments)
|
|
if err := traffic.WriteSpeedFile(*out, csv); err != nil {
|
|
return err
|
|
}
|
|
lines := strings.Count(csv, "\n")
|
|
fmt.Printf("matched: %d speed rows → %s (skipped %d no-location, %d no-policy)\n", lines, *out, skipped, noPolicy)
|
|
return nil
|
|
}
|
|
|
|
func cmdApply(args []string) error {
|
|
fs := flag.NewFlagSet("apply", flag.ExitOnError)
|
|
csv := fs.String("csv", "speeds.csv", "segment-speed CSV")
|
|
data := fs.String("data", "/home/gmp/osm-build/data", "OSRM data dir")
|
|
base := fs.String("base", "northeast", "base dataset prefix (northeast, northeast-us)")
|
|
prefix := fs.String("prefix", "traffic", "suffix for the traffic layer (base-prefix)")
|
|
binary := fs.String("osrm", "/home/gmp/osm-build/build/osrm-customize", "osrm-customize binary")
|
|
skipCopy := fs.Bool("skip-copy", false, "reuse an existing base-prefix dataset")
|
|
fs.Parse(args)
|
|
if _, err := os.Stat(*csv); err != nil {
|
|
return err
|
|
}
|
|
|
|
// OSRM v26's osrm-customize writes the updated contract files IN PLACE
|
|
// next to the input (there is no --output-prefix), so the traffic layer
|
|
// is a renamed copy of the base dataset. CoW copy when the filesystem
|
|
// supports it.
|
|
trafficPrefix := *base + "-" + *prefix
|
|
if !*skipCopy {
|
|
if err := copyDataset(*data, *base, trafficPrefix); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
cmd := exec.Command(*binary, *data+"/"+trafficPrefix, "--segment-speed-file", *csv)
|
|
cmd.Env = append(os.Environ(), "LD_LIBRARY_PATH=/home/gmp/osm-build/prefix/lib")
|
|
fmt.Fprintf(os.Stderr, "running: %s %s %s\n", *binary, *data+"/"+trafficPrefix, *csv)
|
|
start := time.Now()
|
|
out, err := cmd.CombinedOutput()
|
|
if err != nil {
|
|
return fmt.Errorf("osrm-customize: %w\n%s", err, tail(string(out), 1500))
|
|
}
|
|
fmt.Printf("re-customized %s in %s\n", trafficPrefix, time.Since(start).Round(time.Second))
|
|
fmt.Println("serve it: osrm-routed --algorithm mld --port <p> " + *data + "/" + trafficPrefix + ".osrm")
|
|
fmt.Println("provenance: route responses from this layer must carry")
|
|
fmt.Println(" source=computed(traffic), as_of=<fetch time>, policy=assumed")
|
|
return nil
|
|
}
|
|
|
|
// copyDataset copies base.osrm.* → target.osrm.* (reflink when possible).
|
|
func copyDataset(dataDir, base, target string) error {
|
|
targetBase := dataDir + "/" + target + ".osrm"
|
|
if _, err := os.Stat(targetBase + ".ebg"); err == nil {
|
|
fmt.Fprintf(os.Stderr, "%s.* already exists; use --skip-copy to reuse it\n", target)
|
|
return nil
|
|
}
|
|
matches, err := filepath.Glob(dataDir + "/" + base + ".osrm.*")
|
|
if err != nil || len(matches) == 0 {
|
|
return fmt.Errorf("no %s.osrm.* files in %s", base, dataDir)
|
|
}
|
|
start := time.Now()
|
|
for _, m := range matches {
|
|
name := filepath.Base(m)
|
|
dst := dataDir + "/" + target + ".osrm" + strings.TrimPrefix(name, base+".osrm")
|
|
cp := exec.Command("cp", "--reflink=auto", m, dst)
|
|
if err := cp.Run(); err != nil {
|
|
if err2 := exec.Command("cp", m, dst).Run(); err2 != nil {
|
|
return err2
|
|
}
|
|
}
|
|
}
|
|
fmt.Fprintf(os.Stderr, "copied %d files → %s.* (%s)\n", len(matches), target, time.Since(start).Round(time.Second))
|
|
return nil
|
|
}
|
|
|
|
func tail(s string, n int) string {
|
|
if len(s) <= n {
|
|
return s
|
|
}
|
|
return "…" + s[len(s)-n:]
|
|
}
|
|
|
|
func envOr(k, dflt string) string {
|
|
if v := os.Getenv(k); v != "" {
|
|
return v
|
|
}
|
|
return dflt
|
|
}
|