feat(events): add cursor infinite event feed
This commit is contained in:
43
backend/jobs/tests/test_api.py
Normal file
43
backend/jobs/tests/test_api.py
Normal file
@@ -0,0 +1,43 @@
|
||||
from urllib.parse import urlparse
|
||||
|
||||
import pytest
|
||||
from rest_framework.test import APIClient
|
||||
|
||||
from jobs.services import create_job
|
||||
|
||||
|
||||
@pytest.mark.django_db
|
||||
def test_global_events_use_cursor_pagination():
|
||||
create_job(job_type="demo.success")
|
||||
create_job(job_type="demo.fail")
|
||||
create_job(job_type="demo.slow")
|
||||
|
||||
response = APIClient().get("/api/job-events/?limit=2")
|
||||
|
||||
assert response.status_code == 200
|
||||
body = response.json()
|
||||
assert body["next"] is not None
|
||||
assert body["previous"] is None
|
||||
assert len(body["results"]) == 2
|
||||
assert body["results"][0]["id"] > body["results"][1]["id"]
|
||||
|
||||
|
||||
@pytest.mark.django_db
|
||||
def test_global_events_next_cursor_returns_older_events():
|
||||
create_job(job_type="demo.success")
|
||||
create_job(job_type="demo.fail")
|
||||
create_job(job_type="demo.slow")
|
||||
client = APIClient()
|
||||
first_response = client.get("/api/job-events/?limit=2")
|
||||
first_body = first_response.json()
|
||||
parsed_next = urlparse(first_body["next"])
|
||||
next_path = f"{parsed_next.path}?{parsed_next.query}"
|
||||
|
||||
next_response = client.get(next_path)
|
||||
|
||||
assert next_response.status_code == 200
|
||||
body = next_response.json()
|
||||
assert body["next"] is None
|
||||
assert body["previous"] is not None
|
||||
assert len(body["results"]) == 1
|
||||
assert body["results"][0]["id"] < first_body["results"][1]["id"]
|
||||
@@ -1,6 +1,7 @@
|
||||
from django.db import connection
|
||||
from rest_framework import generics, status
|
||||
from rest_framework.exceptions import NotFound, ValidationError
|
||||
from rest_framework.pagination import CursorPagination
|
||||
from rest_framework.response import Response
|
||||
from rest_framework.views import APIView
|
||||
|
||||
@@ -9,6 +10,13 @@ from jobs.serializers import JobCreateSerializer, JobEventSerializer, JobSeriali
|
||||
from jobs.services import create_job, job_stats, retry_failed_job
|
||||
|
||||
|
||||
class JobEventCursorPagination(CursorPagination):
|
||||
page_size = 50
|
||||
page_size_query_param = "limit"
|
||||
max_page_size = 200
|
||||
ordering = "-id"
|
||||
|
||||
|
||||
class HealthAPIView(APIView):
|
||||
def get(self, request):
|
||||
try:
|
||||
@@ -63,21 +71,17 @@ class JobEventsAPIView(APIView):
|
||||
|
||||
class GlobalJobEventsAPIView(APIView):
|
||||
def get(self, request):
|
||||
queryset = JobEvent.objects.order_by("id")
|
||||
after_id = request.query_params.get("after_id")
|
||||
queryset = JobEvent.objects.order_by("-id")
|
||||
event_type = request.query_params.get("type")
|
||||
job_id = request.query_params.get("job_id")
|
||||
if after_id:
|
||||
queryset = queryset.filter(id__gt=after_id)
|
||||
if event_type:
|
||||
queryset = queryset.filter(type=event_type)
|
||||
if job_id:
|
||||
queryset = queryset.filter(job_id=job_id)
|
||||
try:
|
||||
limit = min(int(request.query_params.get("limit", "100")), 500)
|
||||
except ValueError as exc:
|
||||
raise ValidationError({"limit": "Must be an integer."}) from exc
|
||||
return Response(JobEventSerializer(queryset[:limit], many=True).data)
|
||||
|
||||
paginator = JobEventCursorPagination()
|
||||
page = paginator.paginate_queryset(queryset, request, view=self)
|
||||
return paginator.get_paginated_response(JobEventSerializer(page, many=True).data)
|
||||
|
||||
|
||||
class JobStatsAPIView(APIView):
|
||||
|
||||
@@ -1,9 +1,10 @@
|
||||
import type { Health, Job, JobEvent, JobStats } from "./types";
|
||||
import type { CursorPage, Health, Job, JobEvent, JobStats } from "./types";
|
||||
|
||||
const API_BASE_URL = import.meta.env.VITE_API_BASE_URL ?? "http://localhost:8000/api";
|
||||
|
||||
async function request<T>(path: string, init?: RequestInit): Promise<T> {
|
||||
const response = await fetch(`${API_BASE_URL}${path}`, {
|
||||
const url = /^https?:\/\//i.test(path) ? path : `${API_BASE_URL}${path}`;
|
||||
const response = await fetch(url, {
|
||||
headers: { "Content-Type": "application/json", ...(init?.headers ?? {}) },
|
||||
...init
|
||||
});
|
||||
@@ -23,6 +24,13 @@ export type CreateJobBody = {
|
||||
idempotency_key?: string | null;
|
||||
};
|
||||
|
||||
function normalizeCursorPage<T>(value: CursorPage<T> | T[]): CursorPage<T> {
|
||||
if (Array.isArray(value)) {
|
||||
return { next: null, previous: null, results: value };
|
||||
}
|
||||
return value;
|
||||
}
|
||||
|
||||
export const api = {
|
||||
health: () => request<Health>("/health/"),
|
||||
listJobs: () => request<Job[]>("/jobs/"),
|
||||
@@ -31,12 +39,13 @@ export const api = {
|
||||
retryJob: (jobId: string) => request<Job>(`/jobs/${jobId}/retry/`, { method: "POST" }),
|
||||
getStats: () => request<JobStats>("/jobs/stats/"),
|
||||
listJobEvents: (jobId: string) => request<JobEvent[]>(`/jobs/${jobId}/events/`),
|
||||
listEvents: (filters: { after_id?: number; limit?: number; job_id?: string; type?: string } = {}) => {
|
||||
listEvents: (filters: { limit?: number; job_id?: string; type?: string } = {}) => {
|
||||
const params = new URLSearchParams();
|
||||
Object.entries(filters).forEach(([key, value]) => {
|
||||
if (value !== undefined && value !== null && value !== "") params.set(key, String(value));
|
||||
});
|
||||
const query = params.toString();
|
||||
return request<JobEvent[]>(`/job-events/${query ? `?${query}` : ""}`);
|
||||
}
|
||||
return request<CursorPage<JobEvent> | JobEvent[]>(`/job-events/${query ? `?${query}` : ""}`).then(normalizeCursorPage);
|
||||
},
|
||||
listEventsPage: (url: string) => request<CursorPage<JobEvent> | JobEvent[]>(url).then(normalizeCursorPage)
|
||||
};
|
||||
|
||||
@@ -66,7 +66,7 @@ export function DashboardPage() {
|
||||
]);
|
||||
setJobs(jobsData);
|
||||
setStats({ ...emptyStats, ...statsData, by_status: { ...emptyStats.by_status, ...statsData.by_status } });
|
||||
setEvents(eventsData);
|
||||
setEvents(eventsData.results);
|
||||
setHealth(healthData);
|
||||
}, []);
|
||||
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { useCallback, useEffect, useMemo, useState } from "react";
|
||||
import { useCallback, useEffect, useMemo, useRef, useState } from "react";
|
||||
import { Link } from "react-router-dom";
|
||||
import { toast } from "sonner";
|
||||
|
||||
@@ -8,36 +8,71 @@ import { DateTime } from "../components/DateTime";
|
||||
import { EmptyState } from "../components/EmptyState";
|
||||
import type { JobEvent } from "../types";
|
||||
|
||||
const EVENT_PAGE_SIZE = 50;
|
||||
|
||||
function mergeEvents(current: JobEvent[], incoming: JobEvent[]) {
|
||||
const seen = new Set(current.map((event) => event.id));
|
||||
return [...current, ...incoming.filter((event) => !seen.has(event.id))].sort((a, b) => b.id - a.id);
|
||||
}
|
||||
|
||||
export function EventsPage() {
|
||||
const [events, setEvents] = useState<JobEvent[]>([]);
|
||||
const [nextPageUrl, setNextPageUrl] = useState<string | null>(null);
|
||||
const [loading, setLoading] = useState(true);
|
||||
const [loadingMore, setLoadingMore] = useState(false);
|
||||
const [lastRefreshAt, setLastRefreshAt] = useState<string | null>(null);
|
||||
const sentinelRef = useRef<HTMLDivElement | null>(null);
|
||||
|
||||
const refreshInitial = useCallback(async () => {
|
||||
setEvents(await api.listEvents({ limit: 100 }));
|
||||
const refreshFirstPage = useCallback(async (replace = false) => {
|
||||
if (replace) setLoading(true);
|
||||
try {
|
||||
const page = await api.listEvents({ limit: EVENT_PAGE_SIZE });
|
||||
setNextPageUrl((current) => (replace ? page.next : current));
|
||||
setEvents((current) => (replace ? page.results : mergeEvents(current, page.results)));
|
||||
setLastRefreshAt(new Date().toISOString());
|
||||
} finally {
|
||||
if (replace) setLoading(false);
|
||||
}
|
||||
}, []);
|
||||
|
||||
useEffect(() => {
|
||||
void refreshInitial().catch((caught) => toast.error(caught instanceof Error ? caught.message : String(caught)));
|
||||
}, [refreshInitial]);
|
||||
void refreshFirstPage(true).catch((caught) => toast.error(caught instanceof Error ? caught.message : String(caught)));
|
||||
}, [refreshFirstPage]);
|
||||
|
||||
useEffect(() => {
|
||||
const id = window.setInterval(() => {
|
||||
setEvents((current) => {
|
||||
const afterId = current.at(-1)?.id;
|
||||
void api
|
||||
.listEvents({ after_id: afterId, limit: 100 })
|
||||
.then((incoming) => {
|
||||
if (!incoming.length) return;
|
||||
setEvents((latest) => {
|
||||
const seen = new Set(latest.map((event) => event.id));
|
||||
return [...latest, ...incoming.filter((event) => !seen.has(event.id))].slice(-200);
|
||||
});
|
||||
})
|
||||
.catch(() => undefined);
|
||||
return current;
|
||||
});
|
||||
void refreshFirstPage(false).catch(() => undefined);
|
||||
}, 1000);
|
||||
return () => window.clearInterval(id);
|
||||
}, []);
|
||||
}, [refreshFirstPage]);
|
||||
|
||||
const loadMore = useCallback(async () => {
|
||||
if (!nextPageUrl || loadingMore) return;
|
||||
setLoadingMore(true);
|
||||
try {
|
||||
const page = await api.listEventsPage(nextPageUrl);
|
||||
setNextPageUrl(page.next);
|
||||
setEvents((current) => mergeEvents(current, page.results));
|
||||
} catch (caught) {
|
||||
toast.error(caught instanceof Error ? caught.message : String(caught));
|
||||
} finally {
|
||||
setLoadingMore(false);
|
||||
}
|
||||
}, [loadingMore, nextPageUrl]);
|
||||
|
||||
useEffect(() => {
|
||||
const node = sentinelRef.current;
|
||||
if (!node || !nextPageUrl) return undefined;
|
||||
|
||||
const observer = new IntersectionObserver(
|
||||
(entries) => {
|
||||
if (entries.some((entry) => entry.isIntersecting)) void loadMore();
|
||||
},
|
||||
{ rootMargin: "280px 0px" }
|
||||
);
|
||||
observer.observe(node);
|
||||
return () => observer.disconnect();
|
||||
}, [loadMore, nextPageUrl]);
|
||||
|
||||
const newestFirst = useMemo(() => [...events].sort((a, b) => b.id - a.id), [events]);
|
||||
|
||||
@@ -48,9 +83,17 @@ export function EventsPage() {
|
||||
<span className="eyebrow">Audit log</span>
|
||||
<h2>Events</h2>
|
||||
</div>
|
||||
<div className="event-feed-meta">
|
||||
<span className="connection open">Live cursor feed</span>
|
||||
<span className="muted-text">{newestFirst.length} loaded</span>
|
||||
</div>
|
||||
</header>
|
||||
|
||||
<section className="panel">
|
||||
<div className="section-title">
|
||||
<h3>Global Event Feed</h3>
|
||||
<span className="muted-text">Last refresh {lastRefreshAt ? new Date(lastRefreshAt).toLocaleTimeString() : "-"}</span>
|
||||
</div>
|
||||
<div className="table-wrap">
|
||||
<table>
|
||||
<thead>
|
||||
@@ -82,7 +125,17 @@ export function EventsPage() {
|
||||
))}
|
||||
</tbody>
|
||||
</table>
|
||||
{!newestFirst.length && <EmptyState title="No events yet" />}
|
||||
{loading && !newestFirst.length && <EmptyState title="Loading events" />}
|
||||
{!loading && !newestFirst.length && <EmptyState title="No events yet" />}
|
||||
</div>
|
||||
<div className="infinite-sentinel" ref={sentinelRef}>
|
||||
{loadingMore && <span>Loading older events...</span>}
|
||||
{!loadingMore && nextPageUrl && (
|
||||
<button className="secondary-button" type="button" onClick={() => void loadMore()}>
|
||||
Load older events
|
||||
</button>
|
||||
)}
|
||||
{!loadingMore && !nextPageUrl && newestFirst.length > 0 && <span>End of event history</span>}
|
||||
</div>
|
||||
</section>
|
||||
</div>
|
||||
|
||||
@@ -968,6 +968,25 @@ label {
|
||||
justify-content: flex-end;
|
||||
}
|
||||
|
||||
.event-feed-meta {
|
||||
align-items: center;
|
||||
display: flex;
|
||||
flex-wrap: wrap;
|
||||
gap: 8px;
|
||||
justify-content: flex-end;
|
||||
}
|
||||
|
||||
.infinite-sentinel {
|
||||
align-items: center;
|
||||
color: var(--muted);
|
||||
display: flex;
|
||||
font-size: 13px;
|
||||
font-weight: 900;
|
||||
justify-content: center;
|
||||
min-height: 64px;
|
||||
padding-top: 12px;
|
||||
}
|
||||
|
||||
.filter-grid {
|
||||
grid-template-columns: repeat(4, minmax(0, 1fr));
|
||||
}
|
||||
|
||||
@@ -30,6 +30,12 @@ export type JobEvent = {
|
||||
created_at: string;
|
||||
};
|
||||
|
||||
export type CursorPage<T> = {
|
||||
next: string | null;
|
||||
previous: string | null;
|
||||
results: T[];
|
||||
};
|
||||
|
||||
export type JobStats = {
|
||||
total: number;
|
||||
by_status: Record<JobStatus, number>;
|
||||
|
||||
Reference in New Issue
Block a user