later 0.0.41

Distributed Background jobs manager and runner for Rust
<!doctype html>
<html lang="en">
<head>
  <meta charset="utf-8">
  <meta name="viewport" content="width=device-width, initial-scale=1">
  <title>Later job demo</title>
  <style>
    body { font: 16px system-ui, sans-serif; max-width: 760px; margin: 48px auto; padding: 0 20px; color: #1f2937; }
    button, a, input { font: inherit; }
    button, a { display: inline-block; margin: 6px 6px 6px 0; padding: 10px 14px; }
    button { cursor: pointer; }
    form { display: flex; flex-wrap: wrap; align-items: end; gap: 10px; margin: 22px 0; padding: 16px; background: #f3f4f6; border-radius: 8px; }
    label { display: grid; gap: 5px; font-size: 13px; color: #4b5563; }
    input { width: 130px; padding: 9px; }
    button[aria-pressed="true"] { color: white; background: #b42318; border-color: #b42318; }
    .hint { flex-basis: 100%; margin: 0; color: #6b7280; font-size: 13px; }
    #result, #partition-log { min-height: 24px; margin-top: 18px; padding: 12px; background: #f3f4f6; white-space: pre-wrap; font-family: ui-monospace, monospace; font-size: 13px; }
  </style>
</head>
<body>
  <h1>Later job demo</h1>
  <p>Enqueue jobs, then inspect worker activity and 60-second throughput.</p>
  <button data-path="single">Single job</button>
  <button data-path="continuation">Continuation chain</button>
  <button data-path="retry">Fail once, then retry</button>
  <form id="worker-form">
    <strong>Local workers: <output id="worker-count">...</output></strong>
    <button id="remove-worker" type="button">Remove worker</button>
    <button id="add-worker" type="button">Add worker</button>
    <p class="hint">Removing a worker waits for its current job to finish. At least one worker remains.</p>
  </form>
  <form id="bulk-form">
    <label><input id="stress-regular" type="checkbox" checked> Regular jobs</label>
    <label>Jobs <input id="bulk-count" type="number" min="1" max="10000" value="1000" required></label>
    <label>Failures per job <input id="bulk-failures" type="number" min="0" max="10" value="0" required></label>
    <label><input id="stress-sequential" type="checkbox"> Sequential jobs (all partitions)</label>
    <label>Jobs per partition <input id="stress-partition-count" type="number" min="1" max="200" value="10" required></label>
    <button id="bulk-submit" type="submit">Bulk enqueue</button>
    <button id="stress-toggle" type="button" aria-pressed="false">Start stress loop</button>
    <p class="hint">Check regular jobs, sequential jobs (spread across every partition of the
    "orders" topic), or both. The stress loop repeatedly enqueues the checked kinds. Stop waits
    for the current batch to finish.</p>
  </form>
  <h2>Sequential partitions</h2>
  <p>Jobs enqueued with the same key always run in enqueue order; Later hashes the key to a
  partition, the way a Kafka producer key is hashed by the partitioner, so you never pick a
  partition number directly. Jobs with different keys usually land in different partitions and
  run concurrently. Give a key a delay to see its partition not block the others.</p>
  <form id="partition-form">
    <label>Key <input id="partition-key" type="text" maxlength="200" value="customer-1" required></label>
    <label>Jobs <input id="partition-count" type="number" min="1" max="500" value="20" required></label>
    <label>Delay per job (ms) <input id="partition-delay" type="number" min="0" max="5000" value="0" required></label>
    <button id="partition-submit" type="submit">Enqueue by key</button>
    <button id="partition-clear" type="button">Clear log</button>
    <p class="hint">This topic ("orders") has __PARTITION_COUNT__ partitions. Two different
    keys occasionally hash to the same partition; the log below shows exactly where each key
    landed.</p>
  </form>
  <div id="partition-log">No partition jobs yet.</div>
  <a href="/dashboard">Open performance dashboard</a>
  <div id="result">Ready.</div>
  <script>
    const result = document.querySelector("#result");
    const bulkCount = document.querySelector("#bulk-count");
    const bulkFailures = document.querySelector("#bulk-failures");
    const bulkSubmit = document.querySelector("#bulk-submit");
    const stressToggle = document.querySelector("#stress-toggle");
    const stressRegular = document.querySelector("#stress-regular");
    const stressSequential = document.querySelector("#stress-sequential");
    const stressPartitionCount = document.querySelector("#stress-partition-count");
    const workerCount = document.querySelector("#worker-count");
    const addWorker = document.querySelector("#add-worker");
    const removeWorker = document.querySelector("#remove-worker");
    let stressRunning = false;
    let stressStopRequested = false;

    async function post(path) {
      const response = await fetch(path, { method: "POST" });
      const body = await response.json();
      if (!response.ok) throw new Error(body.error || `Request failed (${response.status})`);
      return body;
    }

    async function enqueue(path) {
      result.textContent = "Enqueuing...";
      try {
        const body = await post(path);
        result.textContent = JSON.stringify(body, null, 2);
      } catch (error) {
        result.textContent = error.message;
      }
    }

    async function changeWorkers(path) {
      addWorker.disabled = true;
      removeWorker.disabled = true;
      try {
        const body = await post(path);
        workerCount.textContent = body.workers;
        removeWorker.disabled = body.workers <= 1;
      } catch (error) {
        result.textContent = error.message;
      } finally {
        addWorker.disabled = false;
        if (Number(workerCount.textContent) > 1) removeWorker.disabled = false;
      }
    }

    function updateStressControls() {
      stressToggle.setAttribute("aria-pressed", String(stressRunning));
      stressToggle.textContent = stressStopRequested
        ? "Stopping after current batch..."
        : stressRunning ? "Stop stress loop" : "Start stress loop";
      stressToggle.disabled = stressRunning && stressStopRequested;
      bulkCount.disabled = stressRunning;
      bulkFailures.disabled = stressRunning;
      bulkSubmit.disabled = stressRunning;
      stressRegular.disabled = stressRunning;
      stressSequential.disabled = stressRunning;
      stressPartitionCount.disabled = stressRunning;
    }

    function reportNoJobKindSelected() {
      result.textContent = "Check “Regular jobs”, “Sequential jobs”, or both.";
    }

    // One batch of whichever job kinds are checked. Shared by the one-shot
    // "Bulk enqueue" submit and by each iteration of the stress loop, so
    // both always agree on what "the checked kinds" means.
    async function enqueueOneBatch() {
      let enqueued = 0;
      let failed = 0;
      if (stressRegular.checked) {
        const count = bulkCount.value;
        const failures = bulkFailures.value;
        const body = await post(`/enqueue/bulk?count=${encodeURIComponent(count)}&failures_before_success=${encodeURIComponent(failures)}`);
        enqueued += body.enqueued;
        failed += body.failed;
      }
      if (stressSequential.checked) {
        const jobsPerPartition = stressPartitionCount.value;
        const body = await post(`/enqueue/partition/stress?jobs_per_partition=${encodeURIComponent(jobsPerPartition)}`);
        enqueued += body.enqueued;
      }
      return { enqueued, failed };
    }

    async function runStressLoop() {
      if (!stressRegular.checked && !stressSequential.checked) {
        reportNoJobKindSelected();
        return;
      }
      if (stressRegular.checked && !bulkCount.reportValidity()) return;
      if (stressSequential.checked && !stressPartitionCount.reportValidity()) return;
      stressRunning = true;
      stressStopRequested = false;
      updateStressControls();
      const started = performance.now();
      let batches = 0;
      let enqueued = 0;
      let failed = 0;

      while (!stressStopRequested) {
        try {
          const batch = await enqueueOneBatch();
          batches += 1;
          enqueued += batch.enqueued;
          failed += batch.failed;
          const elapsedSeconds = Math.max((performance.now() - started) / 1000, Number.EPSILON);
          result.textContent = JSON.stringify({
            status: stressStopRequested ? "stopping" : "running",
            batches,
            enqueued,
            failed,
            elapsed_seconds: Number(elapsedSeconds.toFixed(2)),
            average_enqueue_per_second: Number((enqueued / elapsedSeconds).toFixed(2))
          }, null, 2);
        } catch (error) {
          stressStopRequested = true;
          result.textContent = `Stress loop stopped: ${error.message}`;
        }
      }
      stressRunning = false;
      stressStopRequested = false;
      updateStressControls();
    }

    document.querySelectorAll("button[data-path]").forEach(button => {
      button.addEventListener("click", () => enqueue(`/enqueue/${button.dataset.path}`));
    });
    document.querySelector("#bulk-form").addEventListener("submit", async event => {
      event.preventDefault();
      if (!stressRegular.checked && !stressSequential.checked) {
        reportNoJobKindSelected();
        return;
      }
      if (stressRegular.checked && !bulkCount.reportValidity()) return;
      if (stressSequential.checked && !stressPartitionCount.reportValidity()) return;
      result.textContent = "Enqueuing...";
      try {
        const body = await enqueueOneBatch();
        result.textContent = JSON.stringify(body, null, 2);
      } catch (error) {
        result.textContent = error.message;
      }
    });
    stressToggle.addEventListener("click", () => {
      if (stressRunning) {
        stressStopRequested = true;
        updateStressControls();
        return;
      }
      runStressLoop();
    });
    addWorker.addEventListener("click", () => changeWorkers("/workers/add"));
    removeWorker.addEventListener("click", () => changeWorkers("/workers/remove"));
    fetch("/workers")
      .then(response => response.json())
      .then(body => {
        workerCount.textContent = body.workers;
        removeWorker.disabled = body.workers <= 1;
      })
      .catch(error => { result.textContent = error.message; });

    const partitionForm = document.querySelector("#partition-form");
    const partitionKey = document.querySelector("#partition-key");
    const partitionCount = document.querySelector("#partition-count");
    const partitionDelay = document.querySelector("#partition-delay");
    const partitionLog = document.querySelector("#partition-log");
    const partitionClear = document.querySelector("#partition-clear");

    function renderPartitionLog(entries) {
      if (entries.length === 0) {
        partitionLog.textContent = "No partition jobs yet.";
        return;
      }
      const byKey = new Map();
      const byPartition = new Map();
      for (const entry of entries) {
        if (!byKey.has(entry.key)) byKey.set(entry.key, { partition: entry.partition, seqs: [] });
        byKey.get(entry.key).seqs.push(entry.seq);
        if (!byPartition.has(entry.partition)) byPartition.set(entry.partition, new Set());
        byPartition.get(entry.partition).add(entry.key);
      }
      const perKey = [...byKey.entries()]
        .sort((a, b) => a[0].localeCompare(b[0]))
        .map(([key, { partition, seqs }]) => `  "${key}" (partition ${partition}): ${seqs.join(", ")}`)
        .join("\n");
      const perPartition = [...byPartition.entries()]
        .sort((a, b) => a[0] - b[0])
        .map(([partition, keys]) => `  partition ${partition}: ${[...keys].map(k => `"${k}"`).join(", ")}`)
        .join("\n");
      const recent = entries
        .slice(-30)
        .map(entry => `${entry.key}:${entry.seq}`)
        .join("  ");
      partitionLog.textContent =
        `Completion order per key (must be strictly increasing):\n${perKey}\n\n` +
        `Keys grouped by the partition they hashed to (keys sharing a partition share ` +
        `ordering; different partitions run concurrently):\n${perPartition}\n\n` +
        `Last 30 completions in actual finish order (key:seq, interleaving across keys in ` +
        `different partitions is expected and proves they ran concurrently):\n  ${recent}\n\n` +
        `(${entries.length} completed since the log was last cleared)`;
    }

    async function refreshPartitionLog() {
      try {
        const response = await fetch("/partitions/log");
        const body = await response.json();
        renderPartitionLog(body.entries);
      } catch (error) {
        partitionLog.textContent = `Could not load partition log: ${error.message}`;
      }
    }

    partitionForm.addEventListener("submit", async event => {
      event.preventDefault();
      if (!partitionForm.reportValidity()) return;
      const key = partitionKey.value;
      const count = partitionCount.value;
      const delayMs = partitionDelay.value;
      try {
        const body = await post(`/enqueue/partition?key=${encodeURIComponent(key)}&count=${encodeURIComponent(count)}&delay_ms=${encodeURIComponent(delayMs)}`);
        result.textContent = `"${key}" hashed to partition ${body.resolved_partition} (${body.enqueued} jobs enqueued)`;
      } catch (error) {
        partitionLog.textContent = `Enqueue failed: ${error.message}`;
      }
    });
    partitionClear.addEventListener("click", async () => {
      try {
        await post("/partitions/log/clear");
        await refreshPartitionLog();
      } catch (error) {
        partitionLog.textContent = `Clear failed: ${error.message}`;
      }
    });
    refreshPartitionLog();
    setInterval(refreshPartitionLog, 1000);

    window.addEventListener("pagehide", () => { stressStopRequested = true; });
  </script>
</body>
</html>