Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions environments/platform-resharex/consortial.stripes.config.js
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,10 @@ module.exports = {
platformDescription: 'ReShare platform',
hasAllPerms: false,
showDevInfo: true,
reshare: {
showRefresh: true,
liveUpdates: true,
},
staleBundleWarning: { path: '/index.html', header: 'last-modified', interval: 5 },
},
modules: {
Expand Down
4 changes: 4 additions & 0 deletions environments/platform-resharex/stripes.config.js
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,10 @@ module.exports = {
platformDescription: 'ReShare platform',
hasAllPerms: false,
showDevInfo: true,
reshare: {
showRefresh: true,
liveUpdates: true,
},
staleBundleWarning: { path: '/index.html', header: 'last-modified', interval: 5 },
},
modules: {
Expand Down
1 change: 1 addition & 0 deletions platform-rs-dev/stripes.config.js
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ module.exports = {
},
showCost: true,
showRefresh: true,
liveUpdates: true,
showConditions: true,
patronURL: '/users?qindex=barcode&query={patronid}',
},
Expand Down
3 changes: 3 additions & 0 deletions stripes-reshare/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -3,11 +3,14 @@


// Components (grouped with associated helpers)
export { BrokerEventsProvider } from './src/BrokerEvents';
export { default as RequestCacheSync } from './src/RequestCacheSync';
export { default as DirectLink } from './src/DirectLink/DirectLink';
export { default as useCloseDirect } from './src/DirectLink/useCloseDirect';


// Hooks
export { useBrokerEvents } from './src/BrokerEvents';
export { default as useGetSIURL } from './src/useGetSIURL';
export { default as useIntlCallout } from './src/useIntlCallout';
export { default as useIsActionPending } from './src/useIsActionPending';
Expand Down
251 changes: 251 additions & 0 deletions stripes-reshare/src/BrokerEvents.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,251 @@
/**
* One broker event-stream connection for the whole app, and a hook to listen in.
*
* The stream is narrowed server-side to a side and a symbol taken from the
* tenant, so the connection has no per-screen inputs: one per session, surviving
* navigation. Consumers get parsed payloads and decide what to do with them;
* nothing here knows about patron requests. See `RequestCacheSync`.
*
* Inert unless the `reshare.liveUpdates` app-shell flag is set.
*/

import React, { createContext, useContext, useEffect, useLayoutEffect, useRef } from 'react';
import { useStripes } from '@folio/stripes/core';
import useOkapiKy from './useOkapiKy';

// Null when no provider is mounted, which is supported: consumers hear nothing.
const BrokerEventsContext = createContext(null);

// Three missed heartbeats. Keep it a multiple of the broker's own interval,
// `sseHeartbeatInterval` in broker/api/sse_broker.go.
const STALL_TIMEOUT_MS = 45 * 1000;

const BACKOFF_MIN_MS = 1000;
const BACKOFF_MAX_MS = 30 * 1000;

// Half the delay fixed, half jittered, so reconnecting clients spread out.
const backoffMs = (failures) => {
const ceiling = Math.min(BACKOFF_MIN_MS * (2 ** (failures - 1)), BACKOFF_MAX_MS);
return ceiling / 2 + Math.random() * (ceiling / 2);
};

// A missing permission or bad tenant is a standing answer, not a blip.
const isPermanent = (status) => status >= 400 && status < 500 &&
![401, 408, 429].includes(status);

// The broker puts the event name in the JSON body, not in the SSE framing.
const parseFrame = (frame) => {
const data = frame
.split('\n')
.filter(line => line.startsWith('data:'))
.map(line => line.slice('data:'.length).trimStart())
.join('\n');
if (!data) return undefined;
try {
return JSON.parse(data);
} catch (e) {
return undefined;
}
};

const BrokerEventsProvider = ({ side, children }) => {
const stripes = useStripes();
const ky = useOkapiKy();

// Each ky instance pins the token of the render that built it, so the latest
// is read through a ref: token rotation reaches the next reconnect, without
// the current connection restarting on every render.
const kyRef = useRef(ky);
kyRef.current = ky;

const listeners = useRef(new Set());
// Stable identity, so a re-render here does not re-run every consumer's effect.
const subscribe = useRef((listener) => {
listeners.current.add(listener);
return () => { listeners.current.delete(listener); };
}).current;

const enabled = Boolean(stripes.config?.reshare?.liveUpdates && side);
// Tenant is bound into the request headers, so a change of affiliation
// has to reopen the stream.
const tenant = stripes.okapi?.tenant;

useEffect(() => {
if (!enabled) return undefined;

let stopped = false;
let controller = null;
// Distinguishes a watchdog teardown from the request itself failing.
let stalled = false;
// Sets the backoff; delivery clears it.
let failures = 0;
// Set when a connection ends. The stream has no ids and no replay, so
// anything raised before the next one opens is lost, however brief the gap.
let missedEvents = false;

// One listener throwing must not break delivery to the others.
const emit = (method, arg) => listeners.current.forEach((listener) => {
try {
listener[method](arg);
} catch (err) {
// eslint-disable-next-line no-console
console.error(`Broker event listener failed: ${err?.message ?? String(err)}`);
}
});
const emitEvent = (payload) => emit('event', payload);
const emitGap = () => emit('gap');

// Any bytes prove the connection works: reset the retry delay, and tell
// listeners if events were missed while it was down.
const noteDelivery = () => {
if (missedEvents) {
missedEvents = false;
emitGap();
}
failures = 0;
};

const delay = (ms) => new Promise((resolve) => { setTimeout(resolve, ms); });

const streamOnce = async () => {
controller = new AbortController();
stalled = false;
let watchdog;
// Armed before the request, so a response that never arrives at all, such
// as an intermediary buffering the body, times out the same way.
const arm = () => {
clearTimeout(watchdog);
watchdog = setTimeout(() => {
stalled = true;
controller.abort();
}, STALL_TIMEOUT_MS);
};
arm();
let serverClosed = false;
try {
const res = await kyRef.current.get('broker/sse/events', {
searchParams: { side },
// The watchdog, not ky, is what bounds this request.
timeout: false,
signal: controller.signal,
});
const reader = res.body.getReader();
const decoder = new TextDecoder();
let buffer = '';
while (!stopped) {
// eslint-disable-next-line no-await-in-loop
const { done, value } = await reader.read();
if (done) {
serverClosed = true;
break;
}
// Heartbeats count as delivery but yield no events.
arm();
noteDelivery();
buffer += decoder.decode(value, { stream: true });
const frames = buffer.split('\n\n');
buffer = frames.pop();
frames.forEach((frame) => {
const payload = parseFrame(frame);
if (payload) emitEvent(payload);
});
}
} finally {
clearTimeout(watchdog);
}
return serverClosed;
};

const connectLoop = async () => {
while (!stopped) {
try {
// eslint-disable-next-line no-await-in-loop
const serverClosed = await streamOnce();
if (stopped) return;
// A body that closes immediately would otherwise loop tightly.
if (serverClosed) {
failures += 1;
// eslint-disable-next-line no-await-in-loop
await delay(backoffMs(failures));
}
} catch (e) {
if (stopped) return;
if (isPermanent(e.status)) {
// eslint-disable-next-line no-console
console.warn(`Broker event stream refused (${e.status}), giving up: ${e.message}`);
return;
}
// Watchdog aborts are faults and keep their backoff; the only other
// abort is teardown on unmount, which returned above.
failures += 1;
if (failures === 1) {
// eslint-disable-next-line no-console
console.warn(stalled
? `Broker event stream delivered nothing for ${STALL_TIMEOUT_MS}ms, reconnecting`
: `Broker event stream failed, retrying: ${e.message}`);
}
// eslint-disable-next-line no-await-in-loop
await delay(backoffMs(failures));
}

if (stopped) return;
// Reported by the next delivery, not here, so listeners recheck at a
// point where the recheck can succeed.
missedEvents = true;
}
};

connectLoop();

return () => {
stopped = true;
if (controller) controller.abort();
};
}, [enabled, side, tenant]);

return (
<BrokerEventsContext.Provider value={subscribe}>
{children}
</BrokerEventsContext.Provider>
);
};

/**
* Listen to the app's broker event stream.
*
* `onEvent` receives each parsed payload. `onGap` means the stream has not been
* carrying everything, so anything derived from it must be rechecked.
*
* Both are read at call time, so they need not be stable, and both are called
* synchronously: a promise one returns is not awaited. Without a provider above,
* nothing is delivered.
*/
const useBrokerEvents = ({ onEvent, onGap } = {}) => {
const stripes = useStripes();
const subscribe = useContext(BrokerEventsContext);

const handlers = useRef({ onEvent, onGap });
// A layout effect, so the handlers are updated once the render has committed
// and before any frame can arrive. Writing during render would risk callbacks
// from a render React discarded, and a passive effect flushes too late.
useLayoutEffect(() => { handlers.current = { onEvent, onGap }; });

useEffect(() => {
if (!subscribe) return undefined;
return subscribe({
event: (payload) => handlers.current.onEvent?.(payload),
gap: () => handlers.current.onGap?.(),
});
}, [subscribe]);

useEffect(() => {
// Live updates on but no provider above is a wiring mistake, and it looks
// exactly like the feature being switched off.
if (!subscribe && stripes.config?.reshare?.liveUpdates) {
// eslint-disable-next-line no-console
console.warn('useBrokerEvents: no BrokerEventsProvider above this component, so no events will arrive.');
}
}, [subscribe, stripes.config?.reshare?.liveUpdates]);
};

export { BrokerEventsProvider, useBrokerEvents };
Loading
Loading