Jev Events

Sources

Where items come from. Use an integration, wrap any async iterable, accept webhooks or write your own.

A source reads a stream and emits items. Every item has the same small shape, whatever the platform:

interface Item {
  id: string; // unique within its source and connection
  text: string; // what Jev judges
  author?: { id: string; name: string; roles?: string[] };
  at: Date; // when it was created at its origin
  // Lean facts shown to Jev, such as { firstMessage: true }
  facts?: Record<string, JsonValue>;
  raw?: unknown; // the original payload; never sent to Jev
}

Integrations add typed fields on top. An email also has subject, from and labels, a Twitch chat message has channel, firstMessage and bits, and your handlers see them with full types.

Built-in sources

SourcePackageNew itemsGuide
google.gmail.inbox()@jev-events/googleChecked every 15 secondsGmail
google.calendar.invites(), google.calendar.events()@jev-events/googleChecked every 30 secondsGoogle Calendar
slack.messages()@jev-events/slackStreamed, or pushed to your siteSlack
twitch.chat()@jev-events/twitchStreamedTwitch
from(iterable, options?)jev-eventsAs the iterable yields themAnything else
webhook(options?)jev-eventsPushed as HTTP POSTsAnything else
bluesky(options?)jev-events/publicStreamed, with no sign-inAnything else
twitchChat(channel)jev-events/publicStreamed, with no sign-inAnything else

An integration's sources read signed-in accounts. On your machine, that's the account you signed in with npx jev-events auth. In a web app, it's every account your users connected, and the monitor runs once for each of them. See For your users.

How items arrive

A source gets new items in one or more of three ways, which decides where its monitor can run:

DeliveryWhat the source doesWhere it runs
CheckAsks the platform what's new since last time, every everymonitor.start(), a worker, or the cron route
StreamHolds a connection open and emits items as they arrivemonitor.start() or a worker, not a serverless function
PushAnswers the requests the platform posts to your siteThe webhook route in a web app, or its own port

A cursor remembers where a source left off for each connection, such as the newest email it saw, so a restart neither misses nor repeats anything. Cursors are kept in the store: .jev-events/ on your machine, or your database in a web app. An item whose id the monitor already judged for that connection is dropped as a duplicate.

from()

Wrap anything you can iterate: log lines, a queue consumer, rows from a database cursor. Strings become items directly; objects can carry an id, author, at and facts.

any-stream.ts
import { from, monitor, noul, webhook } from "jev-events";

// Any async iterable: log lines, a queue, a database cursor.
const logs = monitor({
  source: from(logLines, { noun: "line" }),
  questions: { outage: noul("Does this log line describe a failure a human should look at?") },
}).on("outage", { min: 0.8 }, (e) => page(e.item.text));

// Or anything that can send a webhook: POST { "text": "…" } to http://127.0.0.1:8787.
const tickets = monitor({
  source: webhook({ port: 8787, secret: process.env.WEBHOOK_SECRET }),
  questions: { angry: noul("Is the customer angry?") },
}).on("angry", { min: 0.7 }, (e) => escalate(e.item.text));

await Promise.all([logs.start(), tickets.start()]);

Prop

Type

A finite iterable ends the source when it runs out, so await logs.run() judges everything and returns the stats. That makes from() handy for backfills and tests.

webhook()

Started on its own, it runs a small HTTP server. POST a string, { "text": "..." }, or an array of them:

curl -X POST http://127.0.0.1:8787/ \
  -H "Authorization: Bearer $WEBHOOK_SECRET" \
  -d '{"text": "Your product broke my build again and nobody answers"}'

It answers 202 with how many items it accepted, 401 without the secret, 400 for a body with no text and 413 for one that's too large. Passed to runtime() in a web app, it opens no port of its own and answers /api/jev/webhook/<monitor id> instead.

Prop

Type

Write your own

A source is a plain object with an id, a platform, and at least one way to get items: check, start or receive. This one checks a status page every 5 minutes, and its cursor keeps a restart from judging an update twice:

status-page.ts
import { fileStore, monitor, noul, type Item, type Source } from "jev-events";

interface Update {
  id: string;
  text: string;
}

/** A status page that lists its updates as JSON, checked for new ones every 5 minutes. */
function statusPage(url: string): Source<Item, "status-page"> {
  return {
    id: `status-page:${url}`,
    platform: "status-page",
    noun: "update",
    defaults: { every: "5m" },
    async check(ctx) {
      const response = await fetch(url, { signal: ctx.signal });
      // A thrown error becomes an "error" event, and the next check runs as usual.
      if (!response.ok) throw new Error(`${url} answered ${response.status}`);
      const updates = (await response.json()) as Update[];

      // The cursor holds what the last check saw. The first check only notes what's there.
      const seen = await ctx.cursor.get<string[]>();
      for (const update of updates) {
        if (seen && !seen.includes(update.id)) await ctx.emit({ id: update.id, text: update.text, at: new Date() });
      }
      await ctx.cursor.set(updates.map((update) => update.id));
    },
  };
}

const vendors = monitor({
  source: statusPage("https://status.example.com/updates.json"),
  questions: { degraded: noul("Does this update say the service is down or degraded now?") },
}).on("degraded", { min: 0.8 }, (e) => console.warn(`Vendor incident: ${e.item.text}`));

// fileStore() keeps the cursor in .jev-events/, so a restart doesn't judge anything twice.
await vendors.start({ store: fileStore() });

Prop

Type

check and start get the same context:

FieldWhat it's for
ctx.emit(item)Hand an item to the monitor. It resolves once the item is queued, not judged
ctx.cursorget() and set(value) where this connection left off. Any JSON value
ctx.signalAborts when the monitor stops. Stop emitting and clean up
ctx.fail(error, { fatal? })Report a problem as an error event. A fatal one also ends this run, and it's started again later
ctx.end()The stream has no more items, so monitor.run() can finish
ctx.connection, ctx.sessionThe account this run reads, and what session() built for it

An error thrown from check becomes an error event, and the next check runs as usual.

Reading signed-in accounts

Set integration, and the monitor runs the source once per connection of that integration. session(ctx) gets the connection with its credentials and builds an API client from them. When it refreshes a token, ctx.saveCredentials() keeps the new one in the store. Throw a SignInError when access was revoked for good: the connection is marked needs-sign-in until its user signs in again, and the monitor emits an error event with needsSignIn: true.

To act on what a source reads, defineAction() turns a platform API call into a native action, with dry-run and protected users built in. The Gmail, Slack and Twitch integrations are built from the same pieces, so their code is a good place to start.

On this page