// 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
# 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 " + *data + "/" + trafficPrefix + ".osrm")
fmt.Println("provenance: route responses from this layer must carry")
fmt.Println(" source=computed(traffic), as_of=, 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
}