import { queryOptions, useMutation, useQuery } from "@octanejs/tanstack-query"; import { useEffect, useRef } from "octane"; import type { TaskRunPage } from "../../../shared/tasks"; import type { CustomerQuerySession } from "../customer-query-session"; import { useCustomerQuerySession } from "../customer-query-provider"; import { callOwner, type SettingsConnection } from "../owner-rpc"; import { activeRun, emptyTaskHistory, retainedRuns, unsettledCursor, } from "../task-history"; function readPage( connection: SettingsConnection, id: string, cursor: TaskRunPage["nextCursor"], signal?: AbortSignal, ) { return callOwner( connection, "listTaskRuns", [id, ...(cursor ? [{ before: cursor }] : [])], signal, ); } export function taskRunsOptions( session: CustomerQuerySession, connection: SettingsConnection, id: string, ) { const queryKey = session.key("tasks", "history", id); return queryOptions({ queryKey, enabled: !!session.getSnapshot().scope && !!connection?.isReady() && !!id, queryFn: async ({ signal }): Promise => { const previous = session.client.getQueryData(queryKey) ?? emptyTaskHistory; const latest = await readPage(connection, id, null, signal); const retained = retainedRuns(latest, previous); const cursor = unsettledCursor(retained.runs, previous); const older = cursor ? await readPage(connection, id, cursor, signal) : null; signal.throwIfAborted(); return { runs: [ ...latest.runs, ...retained.runs.map( (run) => older?.runs.find((item) => item.id === run.id) ?? run, ), ], nextCursor: retained.overlap ? previous.nextCursor : latest.nextCursor, }; }, }); } export function useTaskRunsQuery( connection: SettingsConnection, id: string, busy: boolean, ) { const session = useCustomerQuerySession(); const scope = session.getSnapshot().scope; const options = taskRunsOptions(session, connection, id); const generation = useRef(0); const pagination = useMutation({ mutationFn: async (cursor: NonNullable) => { const token = generation.current; await session.client.cancelQueries({ queryKey: options.queryKey }); if (token !== generation.current) return; const page = await readPage(connection, id, cursor); if (token !== generation.current) return; session.client.setQueryData(options.queryKey, (history) => { if ( !history || history.nextCursor?.id !== cursor.id || history.nextCursor.createdAt !== cursor.createdAt ) return history; return { runs: [ ...history.runs, ...page.runs.filter( (run) => !history.runs.some((item) => item.id === run.id), ), ], nextCursor: page.nextCursor, }; }); }, }); useEffect(() => { pagination.reset(); return () => { generation.current++; }; }, [connection, id, scope, busy]); const query = useQuery({ ...options, enabled: options.enabled && !busy && !pagination.isPending, refetchInterval: (state) => state.state.data?.runs.some(activeRun) ? 3000 : 30_000, }); return { ...query, isFetching: query.isFetching || pagination.isPending, error: query.error ?? pagination.error, more: () => { if ( query.data?.nextCursor && !busy && !query.isFetching && !pagination.isPending && connection?.isReady() ) pagination.mutate(query.data.nextCursor); }, }; }