({
lookbackMs,
clusterId,
}: ClusterFreshnessParams)
| 580 | export const LINE_MAX_COUNT = 10; |
| 581 | |
| 582 | export function useClusterFreshness({ |
| 583 | lookbackMs, |
| 584 | clusterId, |
| 585 | }: ClusterFreshnessParams) { |
| 586 | return useSuspenseQuery({ |
| 587 | queryKey: clusterQueryKeys.clusterFreshness({ |
| 588 | lookbackMs, |
| 589 | clusterId, |
| 590 | }), |
| 591 | queryFn: async ({ queryKey, signal }) => { |
| 592 | const { rows } = await fetchLagHistory({ |
| 593 | params: { |
| 594 | lookback: { |
| 595 | type: "historical", |
| 596 | lookbackMs, |
| 597 | }, |
| 598 | clusterId, |
| 599 | includeSystemObjects: isSystemCluster(clusterId), |
| 600 | }, |
| 601 | requestOptions: { signal }, |
| 602 | queryKey, |
| 603 | }); |
| 604 | |
| 605 | const dataByObjectId = flatGroup(rows, (d) => d.objectId); |
| 606 | |
| 607 | const currentData = dataByObjectId |
| 608 | .map(([_, rowsByObjectId]) => { |
| 609 | // We can assume the last row is the most current because the data is sorted |
| 610 | // by bucket start time. |
| 611 | const lastRow = rowsByObjectId.at(-1); |
| 612 | if (!lastRow) { |
| 613 | return null; |
| 614 | } |
| 615 | return lastRow; |
| 616 | }) |
| 617 | .filter(notNullOrUndefined) |
| 618 | .sort(sortLagInfo) |
| 619 | .slice(0, LINE_MAX_COUNT); |
| 620 | |
| 621 | const dataByBucketStart = group(rows, (d) => d.bucketStart.getTime()); |
| 622 | |
| 623 | const historicalData = [...dataByBucketStart.entries()].map( |
| 624 | ([timestamp, rowsByBucketStart]) => { |
| 625 | const dataPoint: DataPoint = { |
| 626 | timestamp, |
| 627 | lag: {}, |
| 628 | }; |
| 629 | |
| 630 | rowsByBucketStart.forEach((row) => { |
| 631 | dataPoint.lag[row.objectId] = |
| 632 | row.lag !== null |
| 633 | ? { |
| 634 | queryable: true, |
| 635 | totalMs: sumPostgresIntervalMs(row.lag), |
| 636 | interval: row.lag, |
| 637 | schemaName: row.schemaName, |
| 638 | objectName: row.objectName, |
| 639 | } |
no test coverage detected