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
| Source | Package | New items | Guide |
|---|---|---|---|
google.gmail.inbox() | @jev-events/google | Checked every 15 seconds | Gmail |
google.calendar.invites(), google.calendar.events() | @jev-events/google | Checked every 30 seconds | Google Calendar |
slack.messages() | @jev-events/slack | Streamed, or pushed to your site | Slack |
twitch.chat() | @jev-events/twitch | Streamed | Twitch |
from(iterable, options?) | jev-events | As the iterable yields them | Anything else |
webhook(options?) | jev-events | Pushed as HTTP POSTs | Anything else |
bluesky(options?) | jev-events/public | Streamed, with no sign-in | Anything else |
twitchChat(channel) | jev-events/public | Streamed, with no sign-in | Anything 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:
| Delivery | What the source does | Where it runs |
|---|---|---|
| Check | Asks the platform what's new since last time, every every | monitor.start(), a worker, or the cron route |
| Stream | Holds a connection open and emits items as they arrive | monitor.start() or a worker, not a serverless function |
| Push | Answers the requests the platform posts to your site | The 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.
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:
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:
| Field | What it's for |
|---|---|
ctx.emit(item) | Hand an item to the monitor. It resolves once the item is queued, not judged |
ctx.cursor | get() and set(value) where this connection left off. Any JSON value |
ctx.signal | Aborts 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.session | The 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.