@neurowire/ingest
Fetch, detect, and parse (version 0.6.0). Turns a URL into a NeurowireFeed: it fetches with conditional caching and SSRF guards, detects the document kind, parses RSS / Atom / RDF / JSON Feed, auto-detects feeds from HTML, runs the CSS-template engine, and fetches meshes and constructs.
npm install @neurowire/ingestDepends on core, cheerio, fast-xml-parser, and zod.
Fetching feeds
fetchFeed
function fetchFeed(url: string, options?: FetchFeedOptions): Promise<NeurowireFeed>
interface FetchFeedOptions {
template?: FeedTemplate
signal?: AbortSignal
maxDepth?: number
cache?: ConditionalCache
timeoutMs?: number
retries?: number
backoffMs?: number
headers?: Record<string, string>
}Fetch a URL (website, RSS, or Atom) and normalize it to a NeurowireFeed.
| Option | Type | Default | Description |
|---|---|---|---|
template | FeedTemplate? | - | Force a CSS-selector template instead of auto-detecting. |
signal | AbortSignal? | - | Cancel the fetch. |
maxDepth | number? | 3 | Max number of feed-link redirects to follow. |
cache | ConditionalCache? | - | An ETag/Last-Modified response cache, owned by the caller. |
timeoutMs | number? | 15000 | Per-attempt fetch deadline. Set 0 to disable. |
retries | number? | 2 | Max additional attempts after the first. |
backoffMs | number? | 500 | Base delay for exponential backoff with jitter. |
headers | Record<string, string>? | - | Extra request headers for this source. Credential headers stay on the source's origin: a cross-origin redirect or discovered feed link is fetched without them. |
ingestDocument
function ingestDocument(
doc: RawDocument,
options?: FetchFeedOptions,
depth?: number,
): Promise<NeurowireFeed>Turn an already-fetched RawDocument into a feed, without touching the network (useful for testing). Resolution order: explicit template -> discovered feed link -> registry template (by host) -> heuristic auto-detect. Throws when nothing extracts a feed.
Low-level fetch
fetchDocument
function fetchDocument(url: string, options?: FetchOptions): Promise<RawDocument>
interface FetchOptions {
signal?: AbortSignal
cache?: ConditionalCache
validate?: (url: string) => void | Promise<void>
timeoutMs?: number
retries?: number
backoffMs?: number
delay?: (ms: number, signal?: AbortSignal) => Promise<void>
headers?: Record<string, string>
}
interface RawDocument {
url: string
contentType: string
body: string
etag?: string
lastModified?: string
notModified?: boolean
}Fetch a URL over HTTP(S) with a per-attempt timeout and bounded retries. Redirects are followed manually one hop at a time, so the validate guard runs on every hop. Network errors, timeouts, upstream 5xx, and 429 are retried with full-jitter exponential backoff (a 429 honors its Retry-After); a 4xx other than 429, an invalid URL, an SSRF reject, too many redirects, or a caller abort are not retried. A 304 is a success.
| Option | Type | Default | Description |
|---|---|---|---|
signal | AbortSignal? | - | Caller cancellation. |
cache | ConditionalCache? | - | Conditional response cache. |
validate | (url) => void | Promise<void>? | - | Per-URL guard run on the initial URL and every redirect hop. Throw to block (SSRF protection). |
timeoutMs | number? | 15000 | Per-attempt deadline. Set 0 to disable. |
retries | number? | 2 | Max additional attempts. |
backoffMs | number? | 500 | Base backoff delay. |
delay | (ms, signal?) => Promise<void>? | setTimeout-based | Injectable sleep between retries (for tests). |
headers | Record<string, string>? | - | Extra request headers. Override user-agent / accept but never the conditional headers. authorization, proxy-authorization, and cookie are dropped on a cross-origin redirect hop. Not part of the cache key. |
function requestHeaders(
headers: Record<string, string> | undefined,
sameOrigin: boolean,
): Record<string, string>requestHeaders is the per-hop normalizer fetchDocument uses: keys lowercased, conditional headers removed, credential headers removed unless sameOrigin. Exported so a caller that follows links itself can apply the same rule.
RawDocument fields: url (final URL after redirects), contentType, body, etag, lastModified, and notModified (true when served from cache via a 304).
Conditional cache
interface CachedResponse {
url: string
contentType: string
body: string
etag?: string
lastModified?: string
}
interface ConditionalCache {
get(url: string): CachedResponse | undefined
set(url: string, value: CachedResponse): void
}
function createMemoryCache(): ConditionalCache| Export | Description |
|---|---|
CachedResponse | A previously fetched response plus its validators. |
ConditionalCache | A store of cached responses, injected by the caller (the library keeps no global state). |
createMemoryCache() | A simple Map-backed ConditionalCache. |
Detection
type FeedKind = 'nwf' | 'atom' | 'rss' | 'rdf' | 'jsonfeed' | 'html'
function detectKind(contentType: string, body: string): FeedKind| Export | Description |
|---|---|
FeedKind | The detected document kind. |
detectKind(contentType, body) | Classify a document from its Content-Type and body. |
Parsers
Each parser produces a NeurowireFeed from already-parsed input plus a ParseContext.
function parseFeedString(body: string, ctx: ParseContext): NeurowireFeed
function parseNwf(body: string, ctx: ParseContext): NeurowireFeed
function parseAtom(doc: Record<string, unknown>, ctx: ParseContext): NeurowireFeed
function parseRss(doc: Record<string, unknown>, ctx: ParseContext): NeurowireFeed
function parseRdf(doc: Record<string, unknown>, ctx: ParseContext): NeurowireFeed
function parseJsonFeed(raw: unknown, ctx: ParseContext): NeurowireFeed| Export | Description |
|---|---|
parseFeedString(body, ctx) | Sniff a raw string (NWF, JSON Feed, Atom, RSS, or RDF) and dispatch to the right parser. Throws on an unrecognized format. |
parseNwf(body, ctx) | Parse an NWF document back into the model. Throws with the offending line number when the document is malformed, and stamps self from the source URL when the document has none. |
parseAtom(doc, ctx) | Parse a parsed-XML Atom document. |
parseRss(doc, ctx) | Parse a parsed-XML RSS 2.0 document. |
parseRdf(doc, ctx) | Parse a parsed-XML RDF (RSS 1.0) document. |
parseJsonFeed(raw, ctx) | Parse a JSON Feed value. |
HTML auto-detection
function discoverFeedLink($: CheerioAPI, base: string): string | undefined
function autodetect($: CheerioAPI, ctx: ParseContext): NeurowireFeed | null| Export | Description |
|---|---|
discoverFeedLink($, base) | Find a declared feed <link> on the page, resolved against base. Returns undefined when none. |
autodetect($, ctx) | Extract a feed from the page itself (JSON-LD, then semantic HTML). Returns null when nothing extracts. |
The CSS-template engine
A per-site recipe of CSS selectors that turns a listing page into a feed. See taps.
const FeedTemplateSchema: z.ZodType<FeedTemplate>
interface FeedTemplate {
host?: string
feedTitle?: string
item: string
title: string
link?: string
date?: string
summary?: string
author?: string
tags?: string
}
function applyTemplate(
$: CheerioAPI,
template: FeedTemplate,
ctx: ParseContext,
): NeurowireFeed| Field | Type | Description |
|---|---|---|
host | string? | Hostname the template applies to (e.g. blog.example.com). |
feedTitle | string? | Override feed title (otherwise the page <title>). |
item | string | Selector matching each article row. |
title | string | Selector (within an item) for the title text. |
link | string? | Selector for the link (its href). Omit when the item element is itself the link. |
date | string? | Selector for the date; reads [datetime] then text. |
summary | string? | Selector for the summary text. |
author | string? | Selector for the author name. |
tags | string? | Selector matching tag elements. |
| Export | Description |
|---|---|
FeedTemplateSchema | Zod schema for a FeedTemplate. |
applyTemplate($, template, ctx) | Extract a feed from a loaded Cheerio document using a template. |
Proposing a template
interface TemplateProposal {
template: FeedTemplate
matched: number
sampleTitles: string[]
}
function proposeTemplate(html: string, url: string): TemplateProposal | undefined| Export | Description |
|---|---|
TemplateProposal | A guessed template, the count of matched items, and a few sampleTitles. |
proposeTemplate(html, url) | Heuristically propose a template for an HTML page, or undefined when nothing looks like a listing. |
Template registry
function registerTemplate(template: FeedTemplate): void
function findTemplate(url: string): FeedTemplate | undefined
function listTemplates(): FeedTemplate[]| Export | Description |
|---|---|
registerTemplate(template) | Register a per-host template (validated with zod; ignored when it has no host). |
findTemplate(url) | Look up a template for a URL's hostname. |
listTemplates() | All registered templates. |
TIP
@neurowire/taps ships curated templates and registers them with this registry.
Meshes
function fetchMesh(mesh: Mesh, options?: FetchMeshOptions): Promise<NeurowireFeed>
interface FetchMeshOptions {
signal?: AbortSignal
limit?: number
cache?: ConditionalCache
timeoutMs?: number
retries?: number
backoffMs?: number
onSourceError?: (source: { name: string; url: string }, error: unknown) => void
}Fetch every source in a mesh (in parallel) and merge them into one feed. Sources that fail are skipped; throws only if none succeed. A source's own headers (see MeshSource) are passed to its fetch on top of the shared options; other sources never see them, and the default onSourceError warning logs the name, url, and error only.
| Option | Type | Default | Description |
|---|---|---|---|
signal | AbortSignal? | - | Cancellation. |
limit | number? | - | Keep only the newest N merged entries. |
cache | ConditionalCache? | - | A conditional cache shared by every source. |
timeoutMs / retries / backoffMs | number? | 15000 / 2 / 500 | Per-source fetch tuning. |
onSourceError | (source, error) => void? | warns on stderr | Called for each source that failed. A failed source is skipped, never fatal (unless all fail). |
Constructs
type MeshResolver = (ref: string) => Mesh | undefined
interface ConstructPart {
mesh: Mesh
feed: NeurowireFeed
}
interface FetchedConstruct {
name: string
parts: ConstructPart[]
}
function resolveConstructMembers(construct: Construct, resolver?: MeshResolver): Mesh[]
function fetchConstruct(
construct: Construct,
options?: FetchConstructOptions,
): Promise<FetchedConstruct>
function flattenConstruct(
construct: FetchedConstruct,
options?: FlattenConstructOptions,
): NeurowireFeed
interface FetchConstructOptions {
signal?: AbortSignal
limit?: number
cache?: ConditionalCache
resolver?: MeshResolver
concurrency?: number
timeoutMs?: number
retries?: number
backoffMs?: number
onSourceError?: FetchMeshOptions['onSourceError']
onMeshError?: (mesh: Mesh, error: unknown) => void
}
interface FlattenConstructOptions {
limit?: number
}| Export | Description |
|---|---|
MeshResolver | Resolve a mesh reference (by name) to a Mesh, or undefined. The caller supplies the lookup so ingest stays free of filesystem assumptions. |
ConstructPart | One mesh of a fetched construct paired with the feed it merged into. |
FetchedConstruct | A fetched construct: its name plus one merged feed per mesh (grouping preserved). |
resolveConstructMembers(construct, resolver?) | Resolve every member to a concrete Mesh. Inline meshes pass through; { ref } members are looked up. Throws when a ref has no resolver or resolves to nothing. |
fetchConstruct(construct, options?) | Fetch every mesh into its own merged feed with bounded concurrency. Failed meshes are skipped; throws only if none succeed. |
flattenConstruct(construct, options?) | Collapse a fetched construct into one merged feed, tagging entries with their mesh. This is the path the serializers and API use. |
FetchConstructOptions.concurrency defaults to 2: how many meshes to fetch at once. onMeshError is called for each mesh that failed entirely; the default warns on stderr.
Config-backed mesh resolution
interface ConfigResolverOptions {
dirs?: string[]
envVar?: string
subdir?: string
}
function meshConfigDirs(options?: ConfigResolverOptions): string[]
function loadMeshFromConfig(name: string, options?: ConfigResolverOptions): Mesh | undefined
function createConfigMeshResolver(options?: ConfigResolverOptions): MeshResolver
function resolveMeshEnv(mesh: Mesh, env?: NodeJS.ProcessEnv): Mesh
function resolveConstructEnv(construct: Construct, env?: NodeJS.ProcessEnv): Construct
function parseMeshFile(text: string, env?: NodeJS.ProcessEnv): Mesh
function parseConstructFile(text: string, env?: NodeJS.ProcessEnv): Construct| Option | Default | Description |
|---|---|---|
dirs | [] | Extra directories searched before the defaults. |
envVar | NEUROWIRE_MESHES | Env var holding :/,-separated directories. |
subdir | neurowire/meshes | Sub-path under ~/.config (XDG_CONFIG_HOME). |
| Export | Description |
|---|---|
meshConfigDirs(options?) | The directories searched for named mesh files: explicit dirs, then the env var, then ~/.config/neurowire/meshes. |
loadMeshFromConfig(name, options?) | Read a mesh by name (tries <name>.mesh.json then <name>.json). Returns undefined when absent. Rejects path-like names. |
createConfigMeshResolver(options?) | A MeshResolver backed by the config directories. Pass to fetchConstruct({ resolver }). |
resolveMeshEnv(mesh, env?) | Replace ${VAR} in every source's header values from env (default process.env). Throws, naming the mesh, source, header, and variable, when one is unset or empty. Returns a new mesh; the input is untouched. Only for meshes read from trusted local config, never for a remote caller's input. |
resolveConstructEnv(construct, env?) | resolveMeshEnv over a construct's inline meshes; { ref } members pass through. |
parseMeshFile(text, env?) | MeshSchema.parse(JSON.parse(text)) followed by resolveMeshEnv. The one entry point for reading mesh JSON from disk, used by the CLI, the API, and the MCP catalog. |
parseConstructFile(text, env?) | The construct counterpart of parseMeshFile. |
Journal store
The on-disk side of NWFJ: one journal per feed or mesh id, stored as size-capped <id>.<nnnnn>.nwfj segments plus a <id>.manifest.json sidecar. The format itself lives in @neurowire/core.
interface JournalStoreOptions {
dir?: string
maxSegmentBytes?: number
}
function journalConfigDir(): string
function openJournalStore(options?: JournalStoreOptions): JournalStore| Option | Default | Description |
|---|---|---|
dir | journalConfigDir() | Where segments and manifests are written. Created on demand. |
maxSegmentBytes | 5242880 | Rotate to a new segment once the active one reaches this size. |
| Method | Description |
|---|---|
append(id, entries, feed?) | Append entries, dropping ones the journal already holds. Rotates when the active segment is full. Returns { added, skipped, head, segment, rotated }. |
since(id, cursor) | Records after a cursor, plus tooOld when compaction dropped what came before it. |
read(id) | Every record in the journal. |
query(id, query) | Run a JournalQuery, reporting which segment files were scanned and which were skipped. |
head(id) | The journal's head cursor. |
manifest(id) | The segment index, rebuilt automatically when missing or stale. |
verify(id) | Recompute every segment's chain. |
compact(id, keep) | Drop all but the newest keep segments. Returns the removed file names. |
list() | Journal ids present in the directory. |
journalConfigDir() resolves $NEUROWIRE_JOURNAL, else ~/.config/neurowire/journal (honoring XDG_CONFIG_HOME), matching the mesh and tap config directories.
Why queries can skip whole segments
The manifest records each segment's sequence range, date range, and dictionary vocabulary. A query filtering on a tag, author, or source that a segment never declares cannot match anything inside it, so the segment is never opened. The manifest is a cache: delete it and it is rebuilt by rescanning.
Poll engine
The loop behind every live surface: the CLI's tail and --watch, and the API's GET /tail. It is an async generator, so a consumer is a for await and nothing else.
interface PollOptions {
intervalMs?: number
jitter?: number
seen?: Iterable<string>
signal?: AbortSignal
onError?: (error: unknown) => void
delay?: (ms: number, signal?: AbortSignal) => Promise<void>
random?: () => number
}
interface PollTick {
fresh: NeurowireEntry[]
feed: NeurowireFeed
at: number
}
function pollFeed(
load: () => Promise<NeurowireFeed>,
options?: PollOptions,
): AsyncGenerator<PollTick>| Option | Default | Description |
|---|---|---|
intervalMs | 300000 | Base wait between ticks, clamped up to MIN_POLL_INTERVAL_MS (30s). |
jitter | 0.1 | Fraction of the interval added at random, so many pollers do not arrive together. Only ever added, so the interval stays a floor. |
seen | [] | Entry keys already reported, which is how a restart resumes without replaying. |
signal | none | Ends the generator, both between ticks and mid-sleep. |
onError | none | Called when a tick's load throws. The generator survives and waits for the next tick. |
delay | setTimeout | Injectable sleep, so tests never wait on a real timer. |
random | Math.random | Injectable source for the jitter. |
The first tick runs immediately, then the interval applies. Every successful tick is yielded, including ones with no fresh entries, so a caller can report cadence without a second timer. Dedupe uses core's entryKey and newEntries; the engine owns no I/O, so pass a load that already does whatever fetching, filtering, and merging you want. Pair it with a conditional cache so an unchanged source costs a 304 per tick.
See the Tail concept page for what these semantics mean in practice, and CLI tail mode and GET /tail for the surfaces built on them.
| Export | Description |
|---|---|
pollFeed(load, options?) | The generator above. |
resolvePollInterval(ms?) | Apply the default and the floor to a requested interval. |
nextPollDelay(intervalMs, jitter, random) | The jittered wait for one tick, in [interval, interval * (1 + jitter)]. |
MIN_POLL_INTERVAL_MS | 30000. |
DEFAULT_POLL_INTERVAL_MS | 300000. |
DEFAULT_POLL_JITTER | 0.1. |
import { fetchMesh, pollFeed } from '@neurowire/ingest'
const controller = new AbortController()
for await (const { fresh } of pollFeed(() => fetchMesh(mesh), {
intervalMs: 60_000,
signal: controller.signal,
})) {
for (const entry of fresh) console.log(entry.title, entry.link)
}Peer sync
The client half of nwf-sync/1: pull journal deltas from peers you chose to trust, verify each response's hash chain, and append the entries to the local journal store. The server half lives in @neurowire/api.
interface Peer {
url: string
token?: string
journals?: string[]
}
function pullJournal(peer, journalId, store, options?): Promise<PullResult>
function syncPeers(peers, store, options?): Promise<SyncReport>
function listPeerJournals(peer, options?): Promise<string[]>| Export | Description |
|---|---|
pullJournal(peer, id, store, options?) | Pull one journal. Polls the peer's head first and stops there when the recorded cursor matches, otherwise walks segments until the peer reports the transfer complete. Returns { added, skipped, head, remoteHead, requests, bytes, bootstrapped, reset }. |
syncPeers(peers, store, options?) | Pull every selected journal from every peer, collecting failures into a SyncReport rather than throwing: one unreachable peer must not cost you the others' deltas. |
listPeerJournals(peer, options?) | The journal ids a peer publishes. |
SyncError | A sync failure, carrying peer and (when known) journal. |
syncEndpoint(base, name, params?) | Build a /sync/<name> URL against a peer's base URL. |
SYNC_VERSION | The protocol version this client implements. A response naming another version is refused rather than merged. |
| Option | Default | Description |
|---|---|---|
state | the file at peerStatePath() | Where peer cursors are recorded. |
fetch | globalThis.fetch | Injectable, so tests never touch the network. |
signal | - | Caller-driven cancellation. |
timeoutMs | 15000 | Per-request deadline, covering the body read. 0 disables it. |
retries | 2 | Extra attempts on a network failure, a timeout, a 5xx, or a 429. |
backoffMs | 500 | Base for full-jitter exponential backoff. |
delay | a setTimeout sleep | Injectable, so tests need no real waits. |
maxSegments | 500 | Safety valve on how many segments one pull may transfer. |
stats | - | A mutable { requests, bytes } counter, so the cost of a pull survives a throw. |
Peer state
interface PeerStateStore {
get(peer: string, journal: string): JournalCursor | undefined
set(peer: string, journal: string, cursor: JournalCursor): void
entries(): Record<string, JournalCursor>
}
function peerStatePath(): string
function openPeerState(path?: string): PeerStateStore
function createMemoryPeerState(initial?): PeerStateStoreA sequence number is meaningful only inside one journal on one node: when this node appends pulled entries, its own store assigns its own numbering and computes its own chain. So the local head means nothing to the protocol, and what gets recorded is the peer's cursor, keyed by (peer url, journal id) in ~/.config/neurowire/peers-state.json (or $NEUROWIRE_PEERS_STATE).
Written only after the append lands
A crash between the pull and the state write costs one re-pull, which the store's entry-key dedupe absorbs. The reverse order would silently lose entries. Every response is chain- verified before anything from it is appended, so a flipped byte merges nothing.
A peer's journal can be rebuilt, restored, or repointed, leaving the recorded cursor pointing at a journal that no longer exists. A head that moved backwards, or a chain hash that disagrees at the recorded sequence number, is treated as divergence: the cursor resets to zero, the pull starts over, and PullResult.reset says so. Without that check every later sync would report "0 new" forever with a clean exit code.
OPML import
function opmlToMesh(xml: string, name?: string): MeshParse an OPML 2.0 subscription list into a Mesh. Every <outline> with an xmlUrl becomes a source (name = text, else title, else the url's host); nested categories are flattened. The mesh name is name, else the OPML head/title, else "imported". Throws on malformed XML or a schema failure.
TIP
The reverse direction (meshToOpml, constructToOpml) lives in @neurowire/core.
Parse utilities
Lower-level helpers used by the parsers, re-exported for advanced callers.
interface ParseContext {
sourceUrl: string
}
interface FeedDraft {
id?: string
title?: string
home?: string
self?: string
updated?: string
authors?: Person[]
entries: NeurowireEntry[]
}
function resolveUrl(href: string, base: string): string
function normDate(value: string | undefined): string | undefined
function stripHtml(value: string | undefined): string | undefined
function finalizeFeed(draft: FeedDraft, ctx: ParseContext): NeurowireFeed| Export | Description |
|---|---|
ParseContext | The parse context: the sourceUrl the document came from. |
FeedDraft | A loose, in-progress feed that finalizeFeed turns into a valid NeurowireFeed. |
resolveUrl(href, base) | Resolve a possibly-relative href against base. Returns the input on failure. |
normDate(value) | Normalize any parseable date (RFC 822, RFC 3339, ...) to ISO 8601, or undefined. |
stripHtml(value) | Strip tags, decode numeric and common named HTML entities, and collapse whitespace, or undefined when empty. Applied to entry titles and summaries. |
finalizeFeed(draft, ctx) | Fill in defaults, give entries stable ids, and stamp the generator to produce a valid NeurowireFeed. |