Stream delayed live feed events

This commit is contained in:
dirtydishes 2026-05-04 04:59:09 -04:00
parent 85dfebb8f0
commit b88ef2b371
3 changed files with 17 additions and 11 deletions

View file

@ -100,7 +100,7 @@ import {
} from "@islandflow/types";
import { createClient } from "redis";
import { z } from "zod";
import { LiveStateManager, isLiveItemFresh } from "./live";
import { LiveStateManager, shouldFanoutLiveEvent } from "./live";
const service = "api";
const logger = createLogger({ service });
@ -982,14 +982,7 @@ const run = async () => {
) => {
const watermark = await liveState.ingest(ingestChannel, item);
if (
(ingestChannel === "options" ||
ingestChannel === "nbbo" ||
ingestChannel === "equities" ||
ingestChannel === "equity-quotes" ||
ingestChannel === "flow") &&
!isLiveItemFresh(ingestChannel, item)
) {
if (!shouldFanoutLiveEvent(ingestChannel, item)) {
return;
}

View file

@ -289,6 +289,8 @@ export const isLiveItemFresh = (
return now - ts <= thresholdMs;
};
export const shouldFanoutLiveEvent = (_channel: LiveChannel, _item: unknown): boolean => true;
const nextBeforeForItems = <T>(items: T[], cursorOf: (item: T) => Cursor): Cursor | null => {
const last = items.at(-1);
return last ? cursorOf(last) : null;