Skip to content

Built for people who want to own their automations. Join the waitlist for an invite.

Package listing

@kody/fly

src/metrics.ts

228 lines · 7.8 KB · TypeScript
import { clean, positiveInt, requestFlyApi, resolveOrganizationSlug } from './fly-core.ts'
import type { FlyMetricsInput } from './types.ts'

const defaultRateWindow = '5m'
const defaultStep = '60s'

function labelValue(value) {
  return String(value ?? '').replace(/\\/g, '\\\\').replace(/"/g, '\\"').replace(/\n/g, '\\n')
}

function timestampSeconds(value) {
  const cleaned = clean(value)
  if (!cleaned) return null
  if (/^\d+(\.\d+)?$/.test(cleaned)) return cleaned
  const milliseconds = Date.parse(cleaned)
  if (!Number.isFinite(milliseconds)) throw new Error('Invalid timestamp: ' + cleaned)
  return String(milliseconds / 1000)
}

function summarizeSeries(series, maxValues = 5) {
  const metric = series?.metric ?? {}
  if (Array.isArray(series?.value)) {
    return {
      metric,
      value: {
        timestamp: series.value[0],
        value: series.value[1],
      },
    }
  }
  const values = Array.isArray(series?.values) ? series.values : []
  return {
    metric,
    valueCount: values.length,
    first: values[0] ? { timestamp: values[0][0], value: values[0][1] } : null,
    last: values.length ? { timestamp: values.at(-1)[0], value: values.at(-1)[1] } : null,
    sample: values.slice(0, maxValues).map((value) => ({ timestamp: value[0], value: value[1] })),
  }
}

function summarizePrometheusBody(body, { seriesLimit = 10, valueLimit = 5, includeData = false } = {}) {
  const result = Array.isArray(body?.data?.result) ? body.data.result : []
  const summary = {
    status: body?.status ?? null,
    errorType: body?.errorType ?? null,
    error: body?.error ?? null,
    resultType: body?.data?.resultType ?? null,
    resultCount: result.length,
    sample: result.slice(0, seriesLimit).map((series) => summarizeSeries(series, valueLimit)),
  }
  if (includeData) summary.data = body?.data ?? null
  return summary
}

function queryPayload(params) {
  const search = new URLSearchParams()
  search.set('query', clean(params.query))
  if (clean(params.time)) search.set('time', timestampSeconds(params.time))
  return search
}

function queryRangePayload(params) {
  const start = timestampSeconds(params.start)
  const end = timestampSeconds(params.end)
  if (!start || !end) throw new Error('start and end are required for query_range.')
  const search = new URLSearchParams()
  search.set('query', clean(params.query))
  search.set('start', start)
  search.set('end', end)
  search.set('step', clean(params.step) || defaultStep)
  return search
}

async function requestPrometheus(orgSlug, endpoint, payload, auth = {}) {
  return await requestFlyApi(`/prometheus/${encodeURIComponent(orgSlug)}/api/v1/${endpoint}`, {
    account: auth.account,
    secretName: auth.secretName,
    method: 'POST',
    headers: {
      'Content-Type': 'application/x-www-form-urlencoded',
    },
    body: payload,
  })
}

export async function prometheusQuery(params = {}) {
  const query = clean(params.query)
  if (!query) throw new Error('query is required.')
  const orgSlug = await resolveOrganizationSlug(params)
  const body = await requestPrometheus(orgSlug, 'query', queryPayload(params), params)
  return {
    orgSlug,
    query,
    mode: 'query',
    ...summarizePrometheusBody(body, params),
  }
}

export async function prometheusQueryRange(params = {}) {
  const query = clean(params.query)
  if (!query) throw new Error('query is required.')
  const orgSlug = await resolveOrganizationSlug(params)
  const body = await requestPrometheus(orgSlug, 'query_range', queryRangePayload(params), params)
  return {
    orgSlug,
    query,
    mode: 'query_range',
    start: params.start,
    end: params.end,
    step: clean(params.step) || defaultStep,
    ...summarizePrometheusBody(body, params),
  }
}

export function appMetricQueries(params = {}) {
  const appName = clean(params.appName || params.app)
  if (!appName) throw new Error('appName or app is required.')
  const app = labelValue(appName)
  const rateWindow = clean(params.rateWindow) || defaultRateWindow
  return [
    {
      name: 'load-average',
      description: 'Instance load average by instance/minute window.',
      query: `fly_instance_load_average{app="${app}"}`,
    },
    {
      name: 'cpu-non-idle-rate',
      description: 'Non-idle CPU centisecond rate by instance and mode.',
      query: `sum by (app, instance, mode) (rate(fly_instance_cpu{app="${app}",mode!="idle"}[${rateWindow}]))`,
    },
    {
      name: 'cpu-throttle-rate',
      description: 'CPU throttle centisecond rate by instance.',
      query: `sum by (app, instance) (rate(fly_instance_cpu_throttle{app="${app}"}[${rateWindow}]))`,
    },
    {
      name: 'memory-available',
      description: 'Available memory bytes by instance.',
      query: `fly_instance_memory_mem_available{app="${app}"}`,
    },
    {
      name: 'memory-total',
      description: 'Total memory bytes by instance.',
      query: `fly_instance_memory_mem_total{app="${app}"}`,
    },
    {
      name: 'disk-read-sectors-rate',
      description: 'Disk sectors read rate by instance/device.',
      query: `sum by (app, instance, device) (rate(fly_instance_disk_sectors_read{app="${app}"}[${rateWindow}]))`,
    },
    {
      name: 'disk-written-sectors-rate',
      description: 'Disk sectors written rate by instance/device.',
      query: `sum by (app, instance, device) (rate(fly_instance_disk_sectors_written{app="${app}"}[${rateWindow}]))`,
    },
    {
      name: 'volume-used-percent',
      description: 'Fly Volume used percentage by volume id when the app has volumes.',
      query: `fly_volume_used_pct{app="${app}"}`,
    },
    {
      name: 'edge-http-responses',
      description: 'Edge HTTP response count increase by status.',
      query: `sum by (app, status) (increase(fly_edge_http_responses_count{app="${app}"}[${rateWindow}]))`,
    },
  ]
}

export async function runMetricQueries(params = {}) {
  const orgSlug = await resolveOrganizationSlug(params)
  const seriesLimit = positiveInt(params.seriesLimit, 10, 100)
  const valueLimit = positiveInt(params.valueLimit, 5, 100)
  const queries = Array.isArray(params.queries) && params.queries.length
    ? params.queries
    : appMetricQueries(params)
  const useRange = params.range === true || (clean(params.start) && clean(params.end))
  const results = []

  for (const item of queries) {
    const query = typeof item === 'string' ? item : item.query
    if (!clean(query)) continue
    const queryParams = {
      ...params,
      orgSlug,
      query,
      seriesLimit,
      valueLimit,
    }
    const result = useRange
      ? await prometheusQueryRange(queryParams)
      : await prometheusQuery(queryParams)
    results.push({
      name: typeof item === 'string' ? null : item.name ?? null,
      description: typeof item === 'string' ? null : item.description ?? null,
      ...result,
    })
  }

  return {
    orgSlug,
    mode: useRange ? 'query_range' : 'query',
    appName: clean(params.appName || params.app) || null,
    start: useRange ? params.start : null,
    end: useRange ? params.end : null,
    step: useRange ? clean(params.step) || defaultStep : null,
    queryCount: results.length,
    results,
  }
}

/**
 * Query Fly Prometheus metrics for an app or run a custom PromQL query.
 * @param params.appName - App name for bundled CPU/memory/HTTP metric queries.
 * @param params.start - Range query start when querying over a window.
 * @returns Summarized Prometheus series for each bundled or custom query.
 * @example
 * import metrics from 'kody:@kody/fly/metrics'
 * const result = await metrics({ appName: 'my-app', organizationSlug: 'my-org', start: '2026-06-24T14:00:00Z', end: '2026-06-24T15:00:00Z' })
 * // => { orgSlug: 'my-org', queryCount: 6, results: [...] }
 */
export default async function metrics(params: FlyMetricsInput = {}) {
  if (clean(params.query)) {
    return clean(params.start) && clean(params.end)
      ? await prometheusQueryRange(params)
      : await prometheusQuery(params)
  }
  return await runMetricQueries(params)
}