code/day60-capstone-3/frame_pipeline.cuThis is the source used by the lesson and its recorded evidence. Compile commands and expected output live in the directory README.
// SPDX-License-Identifier: MIT
//
// Day 60, capstone 3: a real-time frame pipeline.
//
// 900 synthetic 1920x1080 grayscale frames at a 60 frames-per-second
// cadence: a host producer thread generates each frame into a pinned ring
// buffer, the main thread copies it to one of two device input buffers,
// runs a three-kernel chain (box blur, Sobel magnitude, threshold) that was
// captured once as a CUDA graph, and copies one byte per pixel back. All
// cross-stream ordering is CUDA events; the host never synchronizes inside
// the steady state except to release a ring slot whose copy has finished.
//
// The program measures four things and prints all of them, because the
// grade is per card: the per-frame cost of each component (generate, H2D,
// kernel chain, D2H), the pipeline ceiling (the largest of them), an
// unoverlapped synchronous baseline, and the pipelined rate. The paced run
// then holds the 60 fps cadence and counts dropped frames, and the drop
// counter must read zero.
//
// Every frame is a pure function of its index, so the CPU reference can
// regenerate and check any frame byte for byte. The kernels and the
// reference share the same __host__ __device__ arithmetic, all of it
// integer, so the comparison is memcmp, not a tolerance.
//
// Build: nvcc -std=c++17 -O3 -arch=sm_75 -lineinfo -o frame_pipeline \
// frame_pipeline.cu
// Run: ./frame_pipeline
#include <cmath>
#include <cstdio>
#include <cstdlib>
#include <cstring>
#include <ctime>
#include <atomic>
#include <thread>
#include <vector>
#include <cuda_runtime.h>
#include <nvtx3/nvToolsExt.h>
// The one error macro, byte-identical across the course (day 5).
#define CUDA_CHECK(call) \
do { \
cudaError_t err_ = (call); \
if (err_ != cudaSuccess) { \
std::fprintf(stderr, "CUDA error %s:%d: %s: %s\n", __FILE__, \
__LINE__, #call, cudaGetErrorString(err_)); \
std::exit(EXIT_FAILURE); \
} \
} while (0)
// 1920x1080 is the frame the whole capstone is stated in: one byte per
// pixel, 2,073,600 bytes per frame, 900 frames at 60 fps is a 15 second
// steady state.
constexpr int kRows = 1080;
constexpr int kCols = 1920;
constexpr size_t kFrameBytes = static_cast<size_t>(kRows) * kCols;
// The synthetic source. Frames are cut from a 2048x2048 procedural image by
// a fixed-point pan and zoom, so there is no video file and no decoder, and
// frame t is reproducible on any machine.
constexpr int kBaseDim = 2048;
// Four pinned slots is two more than the pipeline's depth of two, so the
// producer can run ahead of the copy engine without ever touching a slot
// still being read. The static_assert is the safety margin, not taste.
constexpr int kRingSlots = 4;
static_assert(kRingSlots >= 3,
"the pipeline is two frames deep; the ring must be deeper");
constexpr int kCheckFrames = 8; // correctness pass, exercises slot reuse
constexpr int kBaselineFrames = 120;
constexpr int kBurstFrames = 300;
constexpr int kPacedFrames = 900;
constexpr int kPoolFrames = 8; // pregenerated frames the burst producer cycles
constexpr int kTargetFps = 60;
constexpr long kFramePeriodNs = 1000000000L / kTargetFps;
constexpr int kWarmupRuns = 3;
constexpr int kTimedRuns = 10;
constexpr int kThreshold = 96;
// 32x8 = 256 threads, eight warps, and threadIdx.x walks the row so one
// warp's 32 loads and stores are 32 consecutive bytes of one image row.
constexpr int kBlockDimX = 32;
constexpr int kBlockDimY = 8;
constexpr int kGridX = (kCols + kBlockDimX - 1) / kBlockDimX;
constexpr int kGridY = (kRows + kBlockDimY - 1) / kBlockDimY;
static_assert(kBlockDimX * kBlockDimY % 32 == 0,
"block size must be a whole number of warps");
// Producer and consumer speak through three atomics. gProduced is the count
// of frames sitting complete in the ring, gReleased the count whose H2D has
// finished so their slot is reusable, gDropped the deadline misses. The
// producer thread makes no CUDA call; every CUDA call stays on the main
// thread. std::thread and std::atomic arrived on day 59.
static std::atomic<int> gProduced{0};
static std::atomic<int> gReleased{0};
static std::atomic<int> gDropped{0};
// Host clock, CLOCK_MONOTONIC, in milliseconds. Used for two jobs CUDA
// events cannot do: pacing the producer at 60 fps, and timing a pipeline
// that spans three streams plus a host thread. Discipline per read: the
// burst and paced reads sit directly after a cudaDeviceSynchronize, the
// baseline's closing read follows the loop's own cudaStreamSynchronize,
// and the generate reads bracket pure host work on an idle device, so no
// read times device work still in flight.
static double wallMs(void) {
timespec ts;
clock_gettime(CLOCK_MONOTONIC, &ts);
return static_cast<double>(ts.tv_sec) * 1.0e3 +
static_cast<double>(ts.tv_nsec) / 1.0e6;
}
static void sleepNs(long ns) {
if (ns <= 0) {
return;
}
timespec req;
req.tv_sec = ns / 1000000000L;
req.tv_nsec = ns % 1000000000L;
nanosleep(&req, nullptr);
}
// Shared integer arithmetic for the kernel chain and the CPU reference.
// Same source compiled for both sides, so a mismatch is a pipeline bug
// (a stale buffer, a missing event), never a rounding difference.
static __host__ __device__ inline int clampi(int v, int lo, int hi) {
return v < lo ? lo : (v > hi ? hi : v);
}
static __host__ __device__ inline int absi(int v) {
return v < 0 ? -v : v;
}
// 3x3 box blur with edge replication. Integer sum of nine, divided by nine.
static __host__ __device__ inline int blurAt(const unsigned char* in, int row,
int col, int rows, int cols) {
int sum = 0;
for (int p = -1; p <= 1; ++p) {
const int r = clampi(row + p, 0, rows - 1);
for (int q = -1; q <= 1; ++q) {
const int c = clampi(col + q, 0, cols - 1);
sum += in[static_cast<size_t>(r) * cols + c];
}
}
return sum / 9;
}
// Sobel gradient magnitude, |gx| + |gy| scaled into a byte. The true range
// of |gx| + |gy| is 0 to 2040, so >> 3 lands in 0 to 255 exactly.
static __host__ __device__ inline int sobelAt(const unsigned char* in, int row,
int col, int rows, int cols) {
int a[3][3];
for (int p = -1; p <= 1; ++p) {
const int r = clampi(row + p, 0, rows - 1);
for (int q = -1; q <= 1; ++q) {
const int c = clampi(col + q, 0, cols - 1);
a[p + 1][q + 1] = in[static_cast<size_t>(r) * cols + c];
}
}
const int gx =
(a[0][2] + 2 * a[1][2] + a[2][2]) - (a[0][0] + 2 * a[1][0] + a[2][0]);
const int gy =
(a[2][0] + 2 * a[2][1] + a[2][2]) - (a[0][0] + 2 * a[0][1] + a[0][2]);
return (absi(gx) + absi(gy)) >> 3;
}
static __host__ __device__ inline int thresholdAt(int v) {
return v >= kThreshold ? 255 : 0;
}
// One thread owns one pixel of the blurred output.
//
// Memory: threadIdx.x walks a row, so one warp's 32 stores are 32
// consecutive bytes; the nine loads per thread overlap heavily with the
// neighbouring threads' and come out of L1 after the first row.
//
// Launch: any 2D grid covering rows x cols; the guard handles the edge.
__global__ void boxBlur3(const unsigned char* __restrict__ in,
unsigned char* __restrict__ out, int rows, int cols) {
const int col = blockIdx.x * blockDim.x + threadIdx.x;
const int row = blockIdx.y * blockDim.y + threadIdx.y;
if (row < rows && col < cols) {
out[static_cast<size_t>(row) * cols + col] =
static_cast<unsigned char>(blurAt(in, row, col, rows, cols));
}
}
// One thread owns one pixel of the gradient image. Same access shape as
// boxBlur3: contiguous stores, window loads served by L1.
// Launch: any 2D grid covering rows x cols.
__global__ void sobelMag(const unsigned char* __restrict__ in,
unsigned char* __restrict__ out, int rows, int cols) {
const int col = blockIdx.x * blockDim.x + threadIdx.x;
const int row = blockIdx.y * blockDim.y + threadIdx.y;
if (row < rows && col < cols) {
out[static_cast<size_t>(row) * cols + col] =
static_cast<unsigned char>(sobelAt(in, row, col, rows, cols));
}
}
// One thread owns one pixel: 255 where the gradient clears kThreshold, else
// 0. One warp reads and writes 32 consecutive bytes.
// Launch: any 2D grid covering rows x cols.
__global__ void thresholdBinary(const unsigned char* __restrict__ in,
unsigned char* __restrict__ out, int rows,
int cols) {
const int col = blockIdx.x * blockDim.x + threadIdx.x;
const int row = blockIdx.y * blockDim.y + threadIdx.y;
if (row < rows && col < cols) {
const size_t i = static_cast<size_t>(row) * cols + col;
out[i] = static_cast<unsigned char>(thresholdAt(in[i]));
}
}
// The procedural source image: 64-pixel tiles with a diagonal gradient on
// top, so the blur has something to smooth and the Sobel finds an edge
// every 64 pixels.
static void buildBase(unsigned char* h_base) {
for (int row = 0; row < kBaseDim; ++row) {
for (int col = 0; col < kBaseDim; ++col) {
const int tile = ((col & 64) ^ (row & 64)) != 0 ? 190 : 40;
const int grad = ((col + row) >> 5) & 31;
h_base[static_cast<size_t>(row) * kBaseDim + col] =
static_cast<unsigned char>(tile + grad);
}
}
}
// Frame t is a 16.16 fixed-point pan and zoom over the base image, nearest
// neighbour. Pure function of (t, row, col): the pan wraps at 96 pixels and
// the zoom step stays low enough that the widest window (col 1919 at the
// largest step, offset 95) still lands inside the 2048-pixel base.
static void generateFrame(const unsigned char* h_base, unsigned char* h_out,
int t) {
const int offsetX = (t * 2) % 96;
const int offsetY = (t * 3) % 96;
const long step = 54000L + 200L * (t % 64);
for (int row = 0; row < kRows; ++row) {
const long srcY = (static_cast<long>(offsetY) << 16) + row * step;
const unsigned char* h_srcRow =
h_base + (srcY >> 16) * static_cast<long>(kBaseDim);
unsigned char* h_dstRow = h_out + static_cast<size_t>(row) * kCols;
long srcX = static_cast<long>(offsetX) << 16;
for (int col = 0; col < kCols; ++col) {
h_dstRow[col] = h_srcRow[srcX >> 16];
srcX += step;
}
}
}
// CPU reference for the whole chain, same arithmetic as the kernels.
// Obvious loops, no OpenMP; its job is to be right.
static void chainCpu(const unsigned char* h_in, unsigned char* h_tmp,
unsigned char* h_out) {
for (int row = 0; row < kRows; ++row) {
for (int col = 0; col < kCols; ++col) {
h_tmp[static_cast<size_t>(row) * kCols + col] =
static_cast<unsigned char>(
blurAt(h_in, row, col, kRows, kCols));
}
}
for (int row = 0; row < kRows; ++row) {
for (int col = 0; col < kCols; ++col) {
h_out[static_cast<size_t>(row) * kCols + col] =
static_cast<unsigned char>(
sobelAt(h_tmp, row, col, kRows, kCols));
}
}
for (size_t i = 0; i < kFrameBytes; ++i) {
h_out[i] = static_cast<unsigned char>(thresholdAt(h_out[i]));
}
}
// Everything the steady state touches, allocated once in main and freed
// there on every path out.
struct Pipeline {
cudaStream_t copyStream;
cudaStream_t computeStream;
cudaStream_t downloadStream;
// Per parity (frame n uses n & 1): the device buffers and the three
// ordering events. Per ring slot: the event that frees the slot.
unsigned char* d_in[2];
unsigned char* d_blur[2];
unsigned char* d_grad[2];
unsigned char* d_edges[2];
cudaEvent_t copyDone[2];
cudaEvent_t graphDone[2];
cudaEvent_t frameDone[2];
cudaEvent_t slotCopied[kRingSlots];
cudaGraphExec_t chain[2];
unsigned char* h_ring; // pinned, kRingSlots frames
unsigned char* h_out; // pinned, 2 frames
};
// Capture the three-kernel chain once per buffer parity. The graph bakes in
// the pointers, so two parities means two graphs rather than a graph update
// per frame; the launches inside the capture are recorded, not run, which
// is why there is no synchronize before EndCapture.
// snippet: capture
static cudaGraphExec_t captureChain(cudaStream_t stream, unsigned char* in,
unsigned char* blur, unsigned char* grad,
unsigned char* edges) {
const dim3 block(kBlockDimX, kBlockDimY);
const dim3 grid(kGridX, kGridY);
cudaGraph_t graph = nullptr;
CUDA_CHECK(cudaStreamBeginCapture(stream, cudaStreamCaptureModeGlobal));
boxBlur3<<<grid, block, 0, stream>>>(in, blur, kRows, kCols);
CUDA_CHECK(cudaGetLastError());
sobelMag<<<grid, block, 0, stream>>>(blur, grad, kRows, kCols);
CUDA_CHECK(cudaGetLastError());
thresholdBinary<<<grid, block, 0, stream>>>(grad, edges, kRows, kCols);
CUDA_CHECK(cudaGetLastError());
cudaGraphExec_t exec = nullptr;
CUDA_CHECK(cudaStreamEndCapture(stream, &graph));
CUDA_CHECK(cudaGraphInstantiate(&exec, graph, 0));
CUDA_CHECK(cudaGraphDestroy(graph));
return exec;
}
// end snippet
// Issue one frame's work and return without waiting for it. All ordering
// is events on the device side: the copy waits for the graph two frames
// back (input buffer reuse), the graph waits for the copy and for the
// download two frames back (edge buffer reuse), the download waits for the
// graph. A wait on a never-recorded event is a no-op, which is what makes
// frames 0 and 1 legal. The one host wait is the slot release at the end,
// on a copy that the double buffering has already all but finished.
// snippet: issue-frame
static void issueFrame(const Pipeline& p, int n) {
const int buf = n & 1;
const int slot = n % kRingSlots;
nvtxRangePushA("upload");
CUDA_CHECK(cudaStreamWaitEvent(p.copyStream, p.graphDone[buf], 0));
CUDA_CHECK(cudaMemcpyAsync(p.d_in[buf], p.h_ring + slot * kFrameBytes,
kFrameBytes, cudaMemcpyHostToDevice,
p.copyStream));
CUDA_CHECK(cudaEventRecord(p.slotCopied[slot], p.copyStream));
CUDA_CHECK(cudaEventRecord(p.copyDone[buf], p.copyStream));
nvtxRangePop();
nvtxRangePushA("chain");
CUDA_CHECK(cudaStreamWaitEvent(p.computeStream, p.copyDone[buf], 0));
CUDA_CHECK(cudaStreamWaitEvent(p.computeStream, p.frameDone[buf], 0));
CUDA_CHECK(cudaGraphLaunch(p.chain[buf], p.computeStream));
CUDA_CHECK(cudaEventRecord(p.graphDone[buf], p.computeStream));
nvtxRangePop();
nvtxRangePushA("download");
CUDA_CHECK(cudaStreamWaitEvent(p.downloadStream, p.graphDone[buf], 0));
CUDA_CHECK(cudaMemcpyAsync(p.h_out + buf * kFrameBytes, p.d_edges[buf],
kFrameBytes, cudaMemcpyDeviceToHost,
p.downloadStream));
CUDA_CHECK(cudaEventRecord(p.frameDone[buf], p.downloadStream));
nvtxRangePop();
CUDA_CHECK(cudaEventSynchronize(p.slotCopied[slot]));
gReleased.store(n + 1, std::memory_order_release);
}
// end snippet
// The producer thread. Paced mode generates real frames and holds the
// 60 fps deadline; a frame whose ring slot is still occupied at its
// deadline is counted as dropped. A camera would overwrite it; this
// producer counts it and then delivers it anyway, so the last-frame
// correctness check still sees all 900 frames and the gate stays
// dropped == 0 rather than "the check happened to pass". Burst mode
// copies from a pregenerated pool as fast as the ring accepts, so the
// pipeline, not the frame synthesis, is what the burst measures.
// snippet: producer
static void producerLoop(unsigned char* h_ring, const unsigned char* h_base,
const unsigned char* h_pool, int frames, bool paced) {
timespec t0;
clock_gettime(CLOCK_MONOTONIC, &t0);
for (int n = 0; n < frames; ++n) {
if (paced) {
timespec now;
clock_gettime(CLOCK_MONOTONIC, &now);
const long elapsed = (now.tv_sec - t0.tv_sec) * 1000000000L +
(now.tv_nsec - t0.tv_nsec);
sleepNs(static_cast<long>(n) * kFramePeriodNs - elapsed);
if (gReleased.load(std::memory_order_acquire) + kRingSlots <= n) {
gDropped.fetch_add(1, std::memory_order_relaxed);
}
}
while (gReleased.load(std::memory_order_acquire) + kRingSlots <= n) {
sleepNs(50000); // ring full: wait for the copy engine
}
unsigned char* h_slot = h_ring + (n % kRingSlots) * kFrameBytes;
nvtxRangePushA("generate");
if (paced) {
generateFrame(h_base, h_slot, n);
} else {
std::memcpy(h_slot, h_pool + (n % kPoolFrames) * kFrameBytes,
kFrameBytes);
}
nvtxRangePop();
gProduced.store(n + 1, std::memory_order_release);
}
}
// end snippet
// Time one issue function with events on the stream that does the work.
// This is day 9's timeKernel with one change and one reason: the start and
// stop events are recorded on the worker stream, because an event recorded
// on the default stream brackets nothing that runs elsewhere. Warm-up runs
// inside, so every kernel and copy timed here has paid its lazy-loading
// cost before the clock starts.
template <typename IssueFn>
static float timeOnStream(cudaStream_t stream, IssueFn issue) {
cudaEvent_t start, stop;
CUDA_CHECK(cudaEventCreate(&start));
CUDA_CHECK(cudaEventCreate(&stop));
for (int i = 0; i < kWarmupRuns; ++i) {
issue();
}
CUDA_CHECK(cudaDeviceSynchronize());
CUDA_CHECK(cudaGetLastError());
CUDA_CHECK(cudaEventRecord(start, stream));
for (int i = 0; i < kTimedRuns; ++i) {
issue();
}
CUDA_CHECK(cudaEventRecord(stop, stream));
CUDA_CHECK(cudaEventSynchronize(stop));
CUDA_CHECK(cudaGetLastError());
float ms = 0.0f;
CUDA_CHECK(cudaEventElapsedTime(&ms, start, stop));
CUDA_CHECK(cudaEventDestroy(start));
CUDA_CHECK(cudaEventDestroy(stop));
return ms / kTimedRuns;
}
// Compare one delivered frame against its CPU reference. Byte-exact by
// construction; a mismatch is reported with its first index and fails the
// program.
static bool frameMatches(const unsigned char* h_got,
const unsigned char* h_want, int frame) {
if (std::memcmp(h_got, h_want, kFrameBytes) == 0) {
return true;
}
for (size_t i = 0; i < kFrameBytes; ++i) {
if (h_got[i] != h_want[i]) {
std::fprintf(stderr, "frame %d wrong at %zu: got %d, want %d\n",
frame, i, h_got[i], h_want[i]);
break;
}
}
return false;
}
// All phases. Failure paths return EXIT_FAILURE with no resource owned
// here; main owns every allocation and frees it after this returns.
static int runPhases(const Pipeline& p, const unsigned char* h_base,
const unsigned char* h_pool,
const std::vector<unsigned char>& h_refEarly,
const std::vector<unsigned char>& h_refMid,
const std::vector<unsigned char>& h_refLast) {
const double budgetMs = 1000.0 / kTargetFps;
// Phase 1: correctness. Eight frames through the full pipeline,
// generated inline (no producer thread), so slot reuse (frame 4 reuses
// frame 0's slot) and buffer reuse (frame 2 reuses frame 0's device
// buffers) both happen while frames are still in flight. Frames 0, 5
// and 7 are checked byte for byte.
nvtxRangePushA("correctness");
gProduced.store(0, std::memory_order_release);
gReleased.store(0, std::memory_order_release);
for (int n = 0; n < kCheckFrames; ++n) {
generateFrame(h_base, p.h_ring + (n % kRingSlots) * kFrameBytes, n);
issueFrame(p, n);
if (n == 0 || n == 5) {
CUDA_CHECK(cudaEventSynchronize(p.frameDone[n & 1]));
const unsigned char* h_want =
(n == 0) ? h_refEarly.data() : h_refMid.data();
if (!frameMatches(p.h_out + (n & 1) * kFrameBytes, h_want, n)) {
return EXIT_FAILURE;
}
}
}
CUDA_CHECK(cudaDeviceSynchronize());
{
std::vector<unsigned char> h_tmp(kFrameBytes);
std::vector<unsigned char> h_want(kFrameBytes);
std::vector<unsigned char> h_frame(kFrameBytes);
generateFrame(h_base, h_frame.data(), kCheckFrames - 1);
chainCpu(h_frame.data(), h_tmp.data(), h_want.data());
if (!frameMatches(p.h_out + ((kCheckFrames - 1) & 1) * kFrameBytes,
h_want.data(), kCheckFrames - 1)) {
return EXIT_FAILURE;
}
}
std::printf("correctness: frames 0, 5 and %d match the CPU reference\n",
kCheckFrames - 1);
nvtxRangePop();
// Phase 2: the components, one at a time, so the learner can compute
// the ceiling for their own card. Generation is host work, so it is
// timed with the host clock over ten frames; the other three are timed
// with events on their own streams.
nvtxRangePushA("components");
const double genStart = wallMs();
for (int i = 0; i < kTimedRuns; ++i) {
generateFrame(h_base, p.h_ring, i);
}
const double genMs = (wallMs() - genStart) / kTimedRuns;
const float h2dMs = timeOnStream(p.copyStream, [&] {
CUDA_CHECK(cudaMemcpyAsync(p.d_in[0], p.h_ring, kFrameBytes,
cudaMemcpyHostToDevice, p.copyStream));
});
const float chainMs = timeOnStream(p.computeStream, [&] {
CUDA_CHECK(cudaGraphLaunch(p.chain[0], p.computeStream));
});
const float d2hMs = timeOnStream(p.downloadStream, [&] {
CUDA_CHECK(cudaMemcpyAsync(p.h_out, p.d_edges[0], kFrameBytes,
cudaMemcpyDeviceToHost, p.downloadStream));
});
// Two ceilings. The device ceiling is the largest of the three
// device-side components and is what the burst phase can reach, since
// its producer only copies from a pool. The full ceiling adds live
// generation and is what the paced run is up against.
double deviceCeilingMs = h2dMs;
const char* wall = "H2D copy";
if (chainMs > deviceCeilingMs) {
deviceCeilingMs = chainMs;
wall = "kernel chain";
}
if (d2hMs > deviceCeilingMs) {
deviceCeilingMs = d2hMs;
wall = "D2H copy";
}
const double fullCeilingMs =
genMs > deviceCeilingMs ? genMs : deviceCeilingMs;
const char* fullWall = genMs > deviceCeilingMs ? "generate (host)" : wall;
std::printf("components, ms per frame (mean of %d):\n", kTimedRuns);
std::printf(" generate %.3f H2D %.3f chain %.3f D2H %.3f\n", genMs,
h2dMs, chainMs, d2hMs);
std::printf(" device ceiling: %.3f ms per frame (%s) = %.1f fps\n",
deviceCeilingMs, wall, 1000.0 / deviceCeilingMs);
std::printf(
" ceiling with live generation: %.3f ms per frame (%s) "
"= %.1f fps\n",
fullCeilingMs, fullWall, 1000.0 / fullCeilingMs);
std::printf(" budget at %d fps: %.3f ms per frame\n", kTargetFps,
budgetMs);
nvtxRangePop();
// Phase 3: the unoverlapped baseline. Copy, run, copy back, wait, one
// frame at a time on one stream, launched through the same graphs so
// the work is identical. Wall clock, read only after the synchronize
// directly above the read. Frames come from the pool through the same
// pinned ring the pipeline uses, so the two phases pay the same
// staging cost and differ only in overlap.
nvtxRangePushA("baseline");
CUDA_CHECK(cudaDeviceSynchronize());
const double baseStart = wallMs();
for (int n = 0; n < kBaselineFrames; ++n) {
const int buf = n & 1;
unsigned char* h_slot = p.h_ring + (n % kRingSlots) * kFrameBytes;
std::memcpy(h_slot, h_pool + (n % kPoolFrames) * kFrameBytes,
kFrameBytes);
CUDA_CHECK(cudaMemcpyAsync(p.d_in[buf], h_slot, kFrameBytes,
cudaMemcpyHostToDevice, p.computeStream));
CUDA_CHECK(cudaGraphLaunch(p.chain[buf], p.computeStream));
CUDA_CHECK(cudaMemcpyAsync(p.h_out + buf * kFrameBytes, p.d_edges[buf],
kFrameBytes, cudaMemcpyDeviceToHost,
p.computeStream));
CUDA_CHECK(cudaStreamSynchronize(p.computeStream));
}
const double baseMs = (wallMs() - baseStart) / kBaselineFrames;
CUDA_CHECK(cudaGetLastError());
std::printf("baseline, synchronous: %.3f ms per frame = %.1f fps\n", baseMs,
1000.0 / baseMs);
nvtxRangePop();
// Phase 4: the burst. Producer unpaced from the pool, pipeline flat
// out, which is where the timeline shows frame n+1's upload under
// frame n's kernels. The sync before each clock read is explicit.
nvtxRangePushA("burst");
gProduced.store(0, std::memory_order_release);
gReleased.store(0, std::memory_order_release);
CUDA_CHECK(cudaDeviceSynchronize());
const double burstStart = wallMs();
std::thread burstProducer(producerLoop, p.h_ring, h_base, h_pool,
kBurstFrames, false);
for (int n = 0; n < kBurstFrames; ++n) {
while (gProduced.load(std::memory_order_acquire) <= n) {
sleepNs(20000);
}
issueFrame(p, n);
}
burstProducer.join();
CUDA_CHECK(cudaDeviceSynchronize());
const double burstMs = (wallMs() - burstStart) / kBurstFrames;
CUDA_CHECK(cudaGetLastError());
std::printf(
"burst, pipelined: %.3f ms per frame = %.1f fps "
"(%.2f of the device ceiling)\n",
burstMs, 1000.0 / burstMs, deviceCeilingMs / burstMs);
nvtxRangePop();
// Phase 5: the paced run. 900 frames at 60 fps, real synthesis, drop
// counter live. The gate of the whole capstone is at the bottom.
nvtxRangePushA("paced");
gProduced.store(0, std::memory_order_release);
gReleased.store(0, std::memory_order_release);
gDropped.store(0, std::memory_order_release);
CUDA_CHECK(cudaDeviceSynchronize());
const double pacedStart = wallMs();
std::thread pacedProducer(producerLoop, p.h_ring, h_base, h_pool,
kPacedFrames, true);
for (int n = 0; n < kPacedFrames; ++n) {
while (gProduced.load(std::memory_order_acquire) <= n) {
sleepNs(200000);
}
issueFrame(p, n);
}
pacedProducer.join();
CUDA_CHECK(cudaDeviceSynchronize());
const double pacedTotalMs = wallMs() - pacedStart;
CUDA_CHECK(cudaGetLastError());
nvtxRangePop();
const double sustainedFps =
static_cast<double>(kPacedFrames) / (pacedTotalMs / 1000.0);
const int dropped = gDropped.load(std::memory_order_acquire);
std::printf(
"paced: %d frames in %.1f ms, sustained %.2f fps "
"(target %d), dropped %d\n",
kPacedFrames, pacedTotalMs, sustainedFps, kTargetFps, dropped);
if (!frameMatches(p.h_out + ((kPacedFrames - 1) & 1) * kFrameBytes,
h_refLast.data(), kPacedFrames - 1)) {
return EXIT_FAILURE;
}
std::printf("correctness: frame %d matches the CPU reference\n",
kPacedFrames - 1);
if (dropped != 0) {
std::fprintf(stderr,
"FAIL: %d dropped frame(s); the pipeline fell more "
"than %d frames behind its deadline\n",
dropped, kRingSlots);
return EXIT_FAILURE;
}
std::printf("PASS: dropped 0 over %d frames at %d fps\n", kPacedFrames,
kTargetFps);
return EXIT_SUCCESS;
}
int main() {
const int device = 0;
CUDA_CHECK(cudaSetDevice(device));
cudaDeviceProp prop;
CUDA_CHECK(cudaGetDeviceProperties(&prop, device));
std::printf("GPU: %s (compute capability %d.%d, %d copy engine(s))\n",
prop.name, prop.major, prop.minor, prop.asyncEngineCount);
std::printf("frame: %dx%d, %zu bytes; ring: %d pinned slots\n", kCols,
kRows, kFrameBytes, kRingSlots);
// Host-side fixtures: the base image, the burst pool, and the CPU
// reference for every frame the phases check (0 and 5 from the
// correctness pass; the pass's own last frame is rebuilt in place; the
// paced run's final frame).
std::vector<unsigned char> h_base(static_cast<size_t>(kBaseDim) * kBaseDim);
buildBase(h_base.data());
std::vector<unsigned char> h_pool(kPoolFrames * kFrameBytes);
for (int n = 0; n < kPoolFrames; ++n) {
generateFrame(h_base.data(), h_pool.data() + n * kFrameBytes, n);
}
std::vector<unsigned char> h_tmp(kFrameBytes);
std::vector<unsigned char> h_frame(kFrameBytes);
std::vector<unsigned char> h_refEarly(kFrameBytes);
std::vector<unsigned char> h_refMid(kFrameBytes);
std::vector<unsigned char> h_refLast(kFrameBytes);
generateFrame(h_base.data(), h_frame.data(), 0);
chainCpu(h_frame.data(), h_tmp.data(), h_refEarly.data());
generateFrame(h_base.data(), h_frame.data(), 5);
chainCpu(h_frame.data(), h_tmp.data(), h_refMid.data());
generateFrame(h_base.data(), h_frame.data(), kPacedFrames - 1);
chainCpu(h_frame.data(), h_tmp.data(), h_refLast.data());
Pipeline p = {};
CUDA_CHECK(cudaStreamCreate(&p.copyStream));
CUDA_CHECK(cudaStreamCreate(&p.computeStream));
CUDA_CHECK(cudaStreamCreate(&p.downloadStream));
CUDA_CHECK(cudaMallocHost(&p.h_ring, kRingSlots * kFrameBytes));
CUDA_CHECK(cudaMallocHost(&p.h_out, 2 * kFrameBytes));
for (int buf = 0; buf < 2; ++buf) {
CUDA_CHECK(cudaMalloc(&p.d_in[buf], kFrameBytes));
CUDA_CHECK(cudaMalloc(&p.d_blur[buf], kFrameBytes));
CUDA_CHECK(cudaMalloc(&p.d_grad[buf], kFrameBytes));
CUDA_CHECK(cudaMalloc(&p.d_edges[buf], kFrameBytes));
// Ordering events only, so timing is disabled: a timing event
// makes the GPU stamp a clock nobody reads.
CUDA_CHECK(
cudaEventCreateWithFlags(&p.copyDone[buf], cudaEventDisableTiming));
CUDA_CHECK(cudaEventCreateWithFlags(&p.graphDone[buf],
cudaEventDisableTiming));
CUDA_CHECK(cudaEventCreateWithFlags(&p.frameDone[buf],
cudaEventDisableTiming));
}
for (int slot = 0; slot < kRingSlots; ++slot) {
CUDA_CHECK(cudaEventCreateWithFlags(&p.slotCopied[slot],
cudaEventDisableTiming));
}
for (int buf = 0; buf < 2; ++buf) {
p.chain[buf] = captureChain(p.computeStream, p.d_in[buf], p.d_blur[buf],
p.d_grad[buf], p.d_edges[buf]);
}
const int status = runPhases(p, h_base.data(), h_pool.data(), h_refEarly,
h_refMid, h_refLast);
for (int buf = 0; buf < 2; ++buf) {
CUDA_CHECK(cudaGraphExecDestroy(p.chain[buf]));
CUDA_CHECK(cudaEventDestroy(p.copyDone[buf]));
CUDA_CHECK(cudaEventDestroy(p.graphDone[buf]));
CUDA_CHECK(cudaEventDestroy(p.frameDone[buf]));
CUDA_CHECK(cudaFree(p.d_in[buf]));
CUDA_CHECK(cudaFree(p.d_blur[buf]));
CUDA_CHECK(cudaFree(p.d_grad[buf]));
CUDA_CHECK(cudaFree(p.d_edges[buf]));
}
for (int slot = 0; slot < kRingSlots; ++slot) {
CUDA_CHECK(cudaEventDestroy(p.slotCopied[slot]));
}
CUDA_CHECK(cudaFreeHost(p.h_ring));
CUDA_CHECK(cudaFreeHost(p.h_out));
CUDA_CHECK(cudaStreamDestroy(p.copyStream));
CUDA_CHECK(cudaStreamDestroy(p.computeStream));
CUDA_CHECK(cudaStreamDestroy(p.downloadStream));
return status;
}