File size: 8,416 Bytes
9f21d0a
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
import { afterEach, describe, expect, test, vi } from "vitest";

import { parseGpPayload, type GpRecord } from "./gp";
import { laneIndexFor, WorkerSampleSource } from "./sampleSource";
import type { Sgp4Request } from "./sgp4Worker";

// The pool hands each satellite to one worker for the life of the session. Three
// things in WorkerSampleSource depend on that being a pure function of the satnum —
// the per-worker satrec cache, the single `#recordSent` set, and the eviction
// budget — so it is pinned here rather than left to the shape of the hash.

describe("laneIndexFor", () => {
  test("sends a satnum to the same lane every time", () => {
    for (const satnum of ["25544", "00900", "62841", "1", "999999"]) {
      const first = laneIndexFor(satnum, 4);
      for (let repeat = 0; repeat < 100; repeat += 1) {
        expect(laneIndexFor(satnum, 4)).toBe(first);
      }
    }
  });

  test("stays inside the pool", () => {
    for (let laneCount = 1; laneCount <= 8; laneCount += 1) {
      for (let satnum = 0; satnum < 500; satnum += 1) {
        const lane = laneIndexFor(String(satnum), laneCount);
        expect(lane).toBeGreaterThanOrEqual(0);
        expect(lane).toBeLessThan(laneCount);
        expect(Number.isInteger(lane)).toBe(true);
      }
    }
  });

  test("collapses to the single lane", () => {
    // The no-pool case has to stay exactly the old behaviour, including for a
    // degenerate lane count.
    expect(laneIndexFor("25544", 1)).toBe(0);
    expect(laneIndexFor("25544", 0)).toBe(0);
  });

  test("spreads consecutive catalog numbers across the pool", () => {
    // Satellites arrive as a slice of the sorted catalog, so consecutive satnums are
    // the realistic input. A hash that grouped them would leave one worker with the
    // whole activation and the rest idle.
    const laneCount = 4;
    const counts: number[] = Array.from({ length: laneCount }, () => 0);
    const population = 5000;
    for (let satnum = 20_000; satnum < 20_000 + population; satnum += 1) {
      const lane = laneIndexFor(String(satnum), laneCount);
      counts[lane] = (counts[lane] ?? 0) + 1;
    }
    const expected = population / laneCount;
    for (const count of counts) {
      // Within 20% of even. Wide, because this asserts "not degenerate" rather than
      // a property of this particular hash.
      expect(count).toBeGreaterThan(expected * 0.8);
      expect(count).toBeLessThan(expected * 1.2);
    }
  });

  test("spreads non-numeric satnums too", () => {
    // The reason this is not `Number(satnum) % laneCount`: those all parse to NaN,
    // and NaN % n is NaN, so every one of them would land on the same worker.
    const laneCount = 4;
    const seen = new Set<number>();
    for (const satnum of ["ISS", "STARLINK-1007", "2019-074A", "COSMOS 2251 DEB", "T-1", "unknown"]) {
      seen.add(laneIndexFor(satnum, laneCount));
    }
    expect(seen.size).toBeGreaterThan(1);
  });
});

// Giving up is the path that protects the main thread, and the pool multiplied its
// triggers: any one of four workers failing moves the whole catalog to synchronous
// propagation on the thread the workers exist to keep clear. It is also the path
// nothing else exercises — the app only reaches it when something is already wrong.

const TLE = "ISS (ZARYA)\n1 25544U 98067A   18342.69352573  .00002284  00000-0  41838-4 0  9992\n2 25544  51.6407 229.0798 0005166 124.8351 329.3296 15.54069892145658";
const issRecord = (): GpRecord => parseGpPayload(TLE)[0] as GpRecord;
const T0 = Date.UTC(2018, 11, 8);

/** Stands in for a real Worker so the pool can be built and then broken on demand. */
class FakeWorker {
  static instances: FakeWorker[] = [];

  posted: Sgp4Request[] = [];

  terminated = false;

  #listeners = new Map<string, ((event: unknown) => void)[]>();

  constructor() {
    FakeWorker.instances.push(this);
  }

  addEventListener(type: string, handler: (event: unknown) => void): void {
    this.#listeners.set(type, [...(this.#listeners.get(type) ?? []), handler]);
  }

  postMessage(request: Sgp4Request): void {
    this.posted.push(request);
  }

  terminate(): void {
    this.terminated = true;
  }

  emit(type: string, event: unknown): void {
    for (const handler of this.#listeners.get(type) ?? []) handler(event);
  }
}

/** Six cores, so `hardwareConcurrency - 2` gives a pool with more than one lane. */
function poolOfFour(): { source: WorkerSampleSource; errors: ReturnType<typeof vi.spyOn> } {
  FakeWorker.instances = [];
  vi.stubGlobal("navigator", { hardwareConcurrency: 6 });
  vi.stubGlobal("Worker", FakeWorker);
  const errors = vi.spyOn(console, "error").mockImplementation(() => {});
  return { source: new WorkerSampleSource(), errors };
}

const ask = (source: WorkerSampleSource, satnum: string) => source.samplerFor(satnum, issRecord()).samples(T0 - 2_700_000, T0 + 2_700_000);

afterEach(() => {
  vi.unstubAllGlobals();
  vi.restoreAllMocks();
});

describe("WorkerSampleSource pool", () => {
  test("spawns one worker per lane and routes a satnum to the same one every time", async () => {
    const { source } = poolOfFour();
    expect(FakeWorker.instances).toHaveLength(4);

    void ask(source, "25544");
    void ask(source, "25544");
    await Promise.resolve();

    const carrying = FakeWorker.instances.filter((worker) => worker.posted.length > 0);
    expect(carrying).toHaveLength(1);
    expect(carrying[0]?.posted.flatMap((request) => request.commands)).toHaveLength(2);
  });

  test("a single worker failing answers every outstanding request inline", async () => {
    const { source, errors } = poolOfFour();
    // Spread across lanes, and left un-flushed as well as flushed so both the
    // in-flight map and the queue have something in them when the pool dies.
    const flushed = ["25544", "00900", "43013", "62841"].map((satnum) => ask(source, satnum));
    await Promise.resolve();
    const queued = ["11111", "22222"].map((satnum) => ask(source, satnum));

    FakeWorker.instances[0]?.emit("error", { message: "boom" });

    const chunks = await Promise.all([...flushed, ...queued]);
    // Answered, and answered with real samples — the fallback has to be usable,
    // not merely non-hanging.
    expect(chunks).toHaveLength(6);
    for (const chunk of chunks) {
      expect(chunk?.positionsFixed.length).toBeGreaterThan(0);
    }
    expect(source.stats.inlineFallbacks).toBe(6);
    // Every lane goes, not just the one that failed.
    expect(FakeWorker.instances.every((worker) => worker.terminated)).toBe(true);
    // Loud, and once.
    expect(errors).toHaveBeenCalledTimes(1);
  });

  test("requests made after the pool is given up on never reach a worker", async () => {
    const { source } = poolOfFour();
    FakeWorker.instances[0]?.emit("error", { message: "boom" });

    const chunk = await ask(source, "25544");

    expect(chunk?.positionsFixed.length).toBeGreaterThan(0);
    expect(FakeWorker.instances.every((worker) => worker.posted.length === 0)).toBe(true);
    expect(source.stats.inlineFallbacks).toBe(1);
  });

  test("a lane that goes silent while holding work gives up too", async () => {
    vi.useFakeTimers();
    try {
      const { source, errors } = poolOfFour();
      const pending = ask(source, "25544");
      await Promise.resolve();

      vi.advanceTimersByTime(4000);

      expect(await pending).toBeDefined();
      expect(errors).toHaveBeenCalledTimes(1);
      expect(source.stats.inlineFallbacks).toBe(1);
    } finally {
      vi.useRealTimers();
    }
  });

  test("an idle lane is not mistaken for a dead one", async () => {
    vi.useFakeTimers();
    try {
      const { source, errors } = poolOfFour();
      // One lane gets work and answers; the other three hold nothing at all. With a
      // single timer for the pool, their silence would read as death.
      void ask(source, "25544");
      await Promise.resolve();
      const busy = FakeWorker.instances.find((worker) => worker.posted.length > 0);
      const batchId = busy?.posted[0]?.batchId as number;
      busy?.emit("message", { data: { batchId, replies: [{ kind: "unopenable", satnum: "25544", reason: "test" }] } });

      vi.advanceTimersByTime(60_000);

      expect(errors).not.toHaveBeenCalled();
      expect(FakeWorker.instances.some((worker) => worker.terminated)).toBe(false);
    } finally {
      vi.useRealTimers();
    }
  });
});