From 1d63f45af115fda39e1205e328155383fa96c38d Mon Sep 17 00:00:00 2001 From: Carter Francis Date: Wed, 23 Sep 2026 14:26:04 -0500 Subject: [PATCH 1/4] Add get_gpu_frames: zero-copy CuPy access to DE-Server's GPU movie buffers Client.get_gpu_frames() (command 44, GET_GPU_FRAMES) maps the server's two ping-pong movie buffers into this process as CuPy arrays, with no copy, and yields each finished buffer as a batch of images plus per-image frame number, scan position and validity. Leaving the loop body releases the buffer to the server. batches(idle_s) also yields None while waiting, so a consumer can check for stop. The simulated server answers "simulated", and get_gpu_frames returns SimulatedGpuFrames instead: the same interface over NumPy frames of the fake dataset, built in-process at the server's frame rate, so clients can be developed offline. To support that, the fake server: - honours Scan - Repeats and clamps the frame rate (SIMULATED_MAX_FPS), - answers stop_acquisition's UDP stop (it previously hung waiting), and the client records when it started each acquisition. gpu_frames loads kernel32/nvcuda on first use, so it imports anywhere. Claude-Session: https://claude.ai/code/session_017DXgzNPmF7pJt9JMGnKTrT --- deapi/client.py | 38 +++ deapi/gpu_frames.py | 284 ++++++++++++++++++ deapi/simulated_server/fake_server.py | 30 +- deapi/simulated_server/initialize_server.py | 20 ++ .../tests/test_fake_server/test_gpu_frames.py | 62 ++++ 5 files changed, 432 insertions(+), 2 deletions(-) create mode 100644 deapi/gpu_frames.py create mode 100644 deapi/tests/test_fake_server/test_gpu_frames.py diff --git a/deapi/client.py b/deapi/client.py index 966f4d98..25af96d3 100644 --- a/deapi/client.py +++ b/deapi/client.py @@ -81,6 +81,11 @@ def __init__(self): self.commandVersion = commandVersion self.read_only = False self._socket_lock = threading.RLock() + #: How many acquisitions this client has started, and when the last one was sent + #: (time.monotonic). SimulatedGpuFrames times its frames from these: a short + #: simulated acquisition can be over before its consumer first asks. + self.acquisitions_started = 0 + self.acquisition_started_at = 0.0 def set_log_level(self, level): log = logging.getLogger("DECameraClientLib") @@ -1295,6 +1300,7 @@ def start_acquisition( log.debug(" Build Time: %.1f ms", lapsed) step_time = self.GetTime() + started = time.monotonic() response = self._sendCommand(command) if logLevel == logging.DEBUG: lapsed = (self.GetTime() - step_time) * 1000 @@ -1304,6 +1310,9 @@ def start_acquisition( if response: ret = response.acknowledge[0].error != True self.refreshProperties = True + if ret: + self.acquisition_started_at = started + self.acquisitions_started += 1 if logLevel == logging.DEBUG: lapsed = (self.GetTime() - step_time) * 1000 @@ -1989,6 +1998,34 @@ def get_movie_buffer_info(self, movieBufferInfo=None, timeoutMsec=5000): return movieBufferInfo + def get_gpu_frames(self): + """ + Zero-copy CuPy access to the server's GPU movie buffers (Windows, same computer). + + Call before start_acquisition. The returned GpuFrames holds both movie buffers as + CuPy arrays (``.buffers``) and the semaphores; iterating it yields each finished + buffer as ``batch.images`` (cupy float32, images x height x width) and ``batch.info`` + (frame, scanX, scanY, valid per image). The buffer is released back to the server + when the loop advances; the server's processing waits while it is held, so call + ``close()`` on it when done. Requires CuPy. See deapi.gpu_frames. + + Against deapi's simulated server this returns a SimulatedGpuFrames instead: the + same interface, with NumPy batches of the simulated dataset and no CuPy needed. + """ + from deapi.gpu_frames import open_gpu_frames + + return open_gpu_frames(self._attach_gpu_frames(), self) + + def _attach_gpu_frames(self, attach=True): + """Ask DE-Server to hand its GPU movie buffer handles to this process (or detach). + Returns what it answered: [control block handle], or ["simulated"] from the fake + server.""" + command = self._addSingleCommand(self.GET_GPU_FRAMES, None, [str(os.getpid()), attach]) + response = self._sendCommand(command) + if response is False: + raise RuntimeError("DE-Server could not share its GPU movie buffers with this process") + return list(self.__getParameters(response.acknowledge[0])) + def get_movie_buffer( self, movieBuffer, movieBufferSize, numFrames, timeoutMsec=5000 ): @@ -3266,6 +3303,7 @@ def ParseChangedProperties(self, changedProperties, response): LIST_REGISTERS = 40 GET_VIRTUAL_IMAGE_INFO = 41 GET_VIRTUAL_IMAGE = 42 + GET_GPU_FRAMES = 44 MMF_DATA_HEADER_SIZE = 24 diff --git a/deapi/gpu_frames.py b/deapi/gpu_frames.py new file mode 100644 index 00000000..3cba4c39 --- /dev/null +++ b/deapi/gpu_frames.py @@ -0,0 +1,284 @@ +"""Zero-copy CuPy access to DE-Server's GPU movie buffers (Windows, same computer). + + frames = client.get_gpu_frames() # before start_acquisition: maps both buffers + client.start_acquisition(1) + for batch in frames: # one finished movie buffer + batch.images # cupy float32 (images, height, width), server memory + batch.info # numpy view: frame, scanX, scanY, valid per image + # the buffer goes back to the server when the loop advances + +The server processes into two ping-pong movie buffers; while this process holds one, the +server's processing waits for it, so finish (or copy) what you need inside the loop body. +GPU work queued on CuPy's current stream is synchronized before the release; synchronize +work on other streams yourself. + +DEAPI hands this process a handle to a small control block once. The server duplicates the +buffer and semaphore handles into this process itself and lists them there, again whenever +it reallocates the buffers, so no DEAPI round trip is needed during an acquisition. +""" +import ctypes +import time +from ctypes import c_size_t, c_uint64 as u64, wintypes + +import numpy as np + +VERSION, BUFFERS, MAX_IMAGES = 4, 2, 1024 # must match GpuFrameSharing/Protocol.h +H = wintypes.HANDLE +IMAGE = np.dtype([("frame", " picked up here, including queueing), wrap and sync times. + self.profile = dict(contextNs=0, attachNs=0, buffers=0, pickupNs=0, pickupMaxNs=0, + wrapNs=0, wrapMaxNs=0, syncNs=0, syncMaxNs=0) + (self._control_handle,) = attach() + address = _k32.MapViewOfFile(H(self._control_handle), 0xF001F, 0, 0, ctypes.sizeof(_Control)) # all access + if not address: + raise ctypes.WinError() + self._control = _Control.from_address(address) + if self._control.version != VERSION: + raise RuntimeError(f"DE-Server GPU frame protocol {self._control.version}, this client expects {VERSION}") + self.device = self._control.device + start = time.perf_counter_ns() + cp.cuda.Device(self.device).use() + cp.cuda.runtime.free(0) # create CuPy's CUDA context for the driver calls + self.profile["contextNs"] += time.perf_counter_ns() - start + self._load() + + def _count(self, name, ns): + self.profile[name + "Ns"] += ns + self.profile[name + "MaxNs"] = max(self.profile[name + "MaxNs"], ns) + + def _load(self): + """Take the buffer and semaphore handles the server listed (again when they change).""" + cp, listed = self._cp, self._control.client + while True: # 0 while the server updates the list + session = listed.session + filled, free, buffers, detach = listed.filled, list(listed.free), list(listed.buffer), listed.detach + if session and listed.session == session: + break + time.sleep(0.001) + start = time.perf_counter_ns() + self._unload() + self.size = self._control.bufferBytes + self.filled, self.free, self._detach = filled, free, detach + self._pointers = [self._map(handle) for handle in buffers] + # Both movie buffers, flat, before any acquisition; batches are views into these. + self.buffers = [cp.ndarray((self.size // 4,), cp.float32, + cp.cuda.MemoryPointer(cp.cuda.UnownedMemory(p, self.size, self, self.device), 0)) + for p in self._pointers] + self._views, self._consumed, self.session = {}, 0, session + self.profile["attachNs"] += time.perf_counter_ns() - start + + def _map(self, handle): + """Map one server buffer into this process (CUDA VMM import); no data is copied.""" + allocation, pointer, size = u64(), u64(), c_size_t(self.size) + _check(_cu.cuMemImportFromShareableHandle(ctypes.byref(allocation), H(handle), 2)) # Win32 handle + _k32.CloseHandle(H(handle)) + _check(_cu.cuMemAddressReserve(ctypes.byref(pointer), size, c_size_t(0), u64(0), u64(0))) + _check(_cu.cuMemMap(pointer, size, c_size_t(0), allocation, u64(0))) + _check(_cu.cuMemSetAccess(pointer, size, (ctypes.c_int * 3)(1, self.device, 3), c_size_t(1))) # read/write + _cu.cuMemRelease(allocation) # the mapping keeps the memory alive + return pointer.value + + def _unload(self): + if not self._pointers: + return + self._cp.cuda.Device(self.device).synchronize() + self.buffers, self._views = [], {} + for pointer in self._pointers: + _cu.cuMemUnmap(u64(pointer), c_size_t(self.size)) + _cu.cuMemAddressFree(u64(pointer), c_size_t(self.size)) + for handle in (self.filled, *self.free, self._detach): + _k32.CloseHandle(H(handle)) + self._pointers = [] + + def server_profile(self): + """The server's counters (see PROFILE; times in ns).""" + return {name: getattr(self._control, name) for name in PROFILE} + + def __iter__(self): + return self.batches() + + def batches(self, idle_s=None): + """Yield each finished buffer as a Batch; with `idle_s`, also yield None after that + long without one, so the caller can check whether to stop.""" + cp, waited = self._cp, 0.0 + while True: + if _k32.WaitForSingleObject(H(self.filled), 100): # nothing yet + if self._control.client.session != self.session: # server reallocated / restarted + self._load() + waited += 0.1 + if idle_s is not None and waited >= idle_s: + waited = 0.0 + yield None + continue + waited = 0.0 + woke = time.perf_counter_ns() + slot = self._control.slot[self._consumed % BUFFERS] + buffer = slot.buffer + key = (buffer, slot.count, slot.width, slot.height) + if key not in self._views: # arrays over fixed offsets, built once + self._views[key] = cp.ndarray((slot.count, slot.height, slot.width), cp.float32, + self.buffers[buffer].data, strides=(slot.imageBytes, slot.rowBytes, 4)) + info = np.frombuffer(slot.image, IMAGE, slot.count) # view of the shared slot + self._count("pickup", woke - slot.publishedNs) + self._count("wrap", time.perf_counter_ns() - woke) + try: + yield Batch(self._views[key], info) + finally: + synced = time.perf_counter_ns() + cp.cuda.get_current_stream().synchronize() # work queued on the images is done + self._count("sync", time.perf_counter_ns() - synced) + self.profile["buffers"] += 1 + self._consumed += 1 + _k32.ReleaseSemaphore(H(self.free[buffer]), 1, None) + + def close(self): + """Stop sharing: the server gets any held buffer back and stops waiting (no DEAPI call).""" + _k32.SetEvent(H(self._detach)) + self._unload() + _k32.UnmapViewOfFile(ctypes.c_void_p(ctypes.addressof(self._control))) + _k32.CloseHandle(H(self._control_handle)) + + +class SimulatedGpuFrames: + """GpuFrames' interface over deapi's simulated server, for clients developed offline. + + Batches are NumPy float32 frames of the server's own simulated dataset (the TiltGrains + its results come from), built in this process and delivered at the server's frame rate + while it acquires. These are copies and there is no GPU: a stand-in for testing, not + the zero-copy path. `client` reads properties (`client[name]`) and `client.acquiring`, + and carries the Client's `acquisitions_started` / `acquisition_started_at`: the frames + belong to the next acquisition started after this was made, timed from its start. + """ + + BATCH = 32 + + def __init__(self, client): + self._client, self._closed = client, False + self._after = client.acquisitions_started + self.device, self.session = None, 1 + self.profile = dict(buffers=0) + + def server_profile(self): + return {} + + def __iter__(self): + return self.batches() + + def batches(self, idle_s=None): + """As GpuFrames.batches: each batch, and None every `idle_s` seconds without one.""" + from deapi.fake_data.grains import TiltGrains + from deapi.simulated_server.fake_server import SIMULATED_MAX_FPS + + c, waited = self._client, 0.0 + while c.acquisitions_started == self._after: # wait for the acquisition + if self._closed: + return + time.sleep(0.05) + waited += 0.05 + if idle_s is not None and waited >= idle_s: + waited = 0.0 + yield None + scanning = c["Scan - Enable"] == "On" + sx, sy = (int(c["Scan - Size X"]), int(c["Scan - Size Y"])) if scanning else (1, 1) + points = sx * sy + total = points * (int(c["Scan - Repeats"] or 1) if scanning else 1) + k = int(c["Sensor Size X (pixels)"]) + fps = min(float(c["Frames Per Second"]), SIMULATED_MAX_FPS) + data = TiltGrains(x_pixels=sx, y_pixels=sy, kx_pixels=k, ky_pixels=k) + signal = np.asarray(data._signal, dtype=np.float32) + start, done, checked, waited = c.acquisition_started_at, 0, 0.0, 0.0 + while done < total and not self._closed: + now = time.monotonic() + if now - checked > 0.1: # one round trip per 100 ms + checked = now + if not c.acquiring and now < start + total / fps - 0.5: + return # stopped early + due = min(total, int((now - start) * fps)) + if due <= done: + time.sleep(0.005) + waited += 0.005 + if idle_s is not None and waited >= idle_s: + waited = 0.0 + yield None + continue + n = min(self.BATCH, due - done) + frame = np.arange(done, done + n) + x, y = frame % points % sx, frame % points // sx + info = np.zeros(n, IMAGE) + info["frame"], info["scanX"], info["scanY"], info["valid"] = frame, x, y, 1 + yield Batch(signal[data.navigator[x, y]], info) + self.profile["buffers"] += 1 + done += n + waited = 0.0 + + def close(self): + self._closed = True diff --git a/deapi/simulated_server/fake_server.py b/deapi/simulated_server/fake_server.py index 829141e4..04d42771 100644 --- a/deapi/simulated_server/fake_server.py +++ b/deapi/simulated_server/fake_server.py @@ -14,6 +14,11 @@ inp_file = resources.files(deapi) / "prop_dump.json" +#: The fastest this simulated camera reads out, whatever "Frames Per Second" asks for — a +#: real server clamps it to what the readout allows. SimulatedGpuFrames paces its batches +#: by the same number, so the frames last as long as the acquisition does. +SIMULATED_MAX_FPS = 2000.0 + def add_parameter(ack, value): """ @@ -327,6 +332,8 @@ def _respond_to_command(self, command=None): == self.GET_VIRTUAL_IMAGE + commandVersion * 100 ): return self._fake_get_virtual_image(command) + elif command.command[0].command_id == self.GET_GPU_FRAMES + commandVersion * 100: + return self._fake_get_gpu_frames(command) else: raise NotImplementedError( f"Command {command.command[0].command_id} not implemented" @@ -375,6 +382,20 @@ def _fake_set_virtual_mask(self, command): print(f"Virtual mask {mask_id} set with {mask}") return (acknowledge_return,) + def _fake_get_gpu_frames(self, command): + """There is no GPU memory to share: answer "simulated", and the client's + get_gpu_frames builds SimulatedGpuFrames over this server's dataset instead.""" + acknowledge_return = pb.DEPacket() + acknowledge_return.type = pb.DEPacket.P_ACKNOWLEDGE + ack1 = acknowledge_return.acknowledge.add() + ack1.command_id = command.command[0].command_id + add_parameter(ack1, "simulated") + return (acknowledge_return,) + + def stop(self): + """End the acquisition now (the client's UDP stop, see initialize_server).""" + self.end_time = min(self.end_time, time.time()) + def _fake_list_cameras(self, command): acknowledge_return = pb.DEPacket() acknowledge_return.type = pb.DEPacket.P_ACKNOWLEDGE @@ -403,7 +424,11 @@ def _fake_start_acquisition(self, command): acknowledge_return = pb.DEPacket() num_acq = command.command[0].parameter[0].p_int if self["Scan - Enable"] == "On": - frames = int(self["Scan - Size X"]) * int(self["Scan - Size Y"]) + frames = ( + int(self["Scan - Size X"]) + * int(self["Scan - Size Y"]) + * int(self["Scan - Repeats"] or 1) + ) self._initialize_data( int(self["Scan - Size X"]), int(self["Scan - Size Y"]), @@ -419,7 +444,7 @@ def _fake_start_acquisition(self, command): ) frames = num_acq self.number_of_frames_requested = frames - fps = float(self["Frames Per Second"]) + fps = min(float(self["Frames Per Second"]), SIMULATED_MAX_FPS) total_time = frames * num_acq / fps self.start_time = time.time() print(f"Acquisition started for {total_time} seconds") @@ -921,3 +946,4 @@ def _fake_get_virtual_image(self, command): SET_CLIENT_READ_ONLY = 31 GET_VIRTUAL_IMAGE_INFO = 41 GET_VIRTUAL_IMAGE = 42 + GET_GPU_FRAMES = 44 diff --git a/deapi/simulated_server/initialize_server.py b/deapi/simulated_server/initialize_server.py index ce533217..1e797e80 100644 --- a/deapi/simulated_server/initialize_server.py +++ b/deapi/simulated_server/initialize_server.py @@ -1,4 +1,5 @@ import sys +import threading import traceback from deapi.simulated_server.fake_server import FakeServer @@ -28,6 +29,22 @@ def _recv_exact(conn, n): return buf +def _serve_stop(host, port, current): + """The client's stop_acquisition: a UDP "PyClientStopAcq" to the same port, answered + "Stopped". Without it a stop against this server would wait for a reply forever.""" + with socket.socket(socket.AF_INET, socket.SOCK_DGRAM) as udp: + try: + udp.bind((host, port)) + except OSError as e: + sys.stderr.write(f"stop_acquisition will not be answered: UDP {port}: {e}\n") + return + while True: + message, addr = udp.recvfrom(64) + if message.startswith(b"PyClientStopAcq") and current: + current[-1].stop() + udp.sendto(b"Stopped", addr) + + # Defining main function def main(port=13240): parser = argparse.ArgumentParser() @@ -52,6 +69,8 @@ def main(port=13240): f" Port: {PORT} \n" ) sys.stderr.flush() + current = [] # the server of the connection being served, for the stop listener + threading.Thread(target=_serve_stop, args=(HOST, PORT, current), daemon=True).start() while True: conn, addr = server_socket.accept() # What waits for a connection # Guard against a stalled client leaving _recv_exact blocked forever. @@ -60,6 +79,7 @@ def main(port=13240): # somehow delivers a short write. conn.settimeout(120) server = FakeServer(socket=conn) + current[:] = [server] connected = True while connected: try: diff --git a/deapi/tests/test_fake_server/test_gpu_frames.py b/deapi/tests/test_fake_server/test_gpu_frames.py new file mode 100644 index 00000000..618ba07d --- /dev/null +++ b/deapi/tests/test_fake_server/test_gpu_frames.py @@ -0,0 +1,62 @@ +"""get_gpu_frames against the simulated server: the same interface as the real, zero-copy +GpuFrames, with NumPy batches of the simulated dataset.""" + +import time + +import numpy as np +import pytest + +from deapi.gpu_frames import SimulatedGpuFrames + + +@pytest.fixture(autouse=True) +def _restore(client): + """The client is shared by the whole session: leave the scan as it was found.""" + before = {name: client[name] for name in ("Scan - Enable", "Scan - Size X", "Scan - Size Y", + "Scan - Repeats", "Frames Per Second")} + yield + for name, value in before.items(): + client[name] = value + + +def _scan(client, size, repeats): + client["Scan - Enable"] = "On" + client["Scan - Size X"] = size + client["Scan - Size Y"] = size + client["Scan - Repeats"] = repeats + client["Frames Per Second"] = 1000 + + +def test_simulated_frames_cover_every_scan_position(client): + _scan(client, 8, 2) + frames = client.get_gpu_frames() + assert isinstance(frames, SimulatedGpuFrames) + client.start_acquisition(1) + seen, positions, deadline = 0, set(), time.monotonic() + 30 + for batch in frames.batches(idle_s=0.2): + assert time.monotonic() < deadline + if batch is None: + continue + assert batch.images.dtype == np.float32 and batch.images.ndim == 3 + assert batch.info["valid"].all() + seen += len(batch.images) + positions.update(zip(batch.info["scanX"].tolist(), batch.info["scanY"].tolist())) + frames.close() + assert seen == 8 * 8 * 2 + assert positions == {(x, y) for x in range(8) for y in range(8)} + + +def test_stop_ends_the_simulated_frames(client): + _scan(client, 8, 1000) + frames = client.get_gpu_frames() + client.start_acquisition(1) + seen, started, deadline = 0, time.monotonic(), time.monotonic() + 30 + for batch in frames.batches(idle_s=0.2): + assert time.monotonic() < deadline + if batch is not None: + seen += len(batch.images) + if seen and time.monotonic() - started > 0.5 and client.acquiring: + assert client.stop_acquisition() + frames.close() + assert 0 < seen < 8 * 8 * 1000 + assert not client.acquiring From 45ebcbfcd61db34d230b1f90883d46e287dfb27e Mon Sep 17 00:00:00 2001 From: Carter Francis Date: Wed, 23 Sep 2026 14:27:11 -0500 Subject: [PATCH 2/4] Simulated server: answer external images (a HAADF detector) external_image1-4 were unsupported, and an unsupported result drops the connection. They are now the annular dark-field image of the fake dataset, which is what an external HAADF detector sees. Claude-Session: https://claude.ai/code/session_017DXgzNPmF7pJt9JMGnKTrT --- deapi/simulated_server/fake_server.py | 7 +++++++ deapi/tests/test_fake_server/test_gpu_frames.py | 10 ++++++++++ 2 files changed, 17 insertions(+) diff --git a/deapi/simulated_server/fake_server.py b/deapi/simulated_server/fake_server.py index 04d42771..0f708170 100644 --- a/deapi/simulated_server/fake_server.py +++ b/deapi/simulated_server/fake_server.py @@ -732,6 +732,13 @@ def _fake_get_result(self, command): ] image = self.fake_data.get_virtual_image(image, method=calculation_type) image = image.astype(pixel_format_dict[pixel_format]) + elif 22 <= frame_type < 26: # external images: an annular (HAADF) detector + ky, kx = self.fake_data.signal.shape[1:] + yy, xx = np.mgrid[0:ky, 0:kx] + r = np.hypot(yy - ky / 2, xx - kx / 2) + mask = np.where(r > 0.15 * min(kx, ky), 2, 1).astype(np.int8) + image = self.fake_data.get_virtual_image(mask, method="Sum") + image = image.astype(pixel_format_dict[pixel_format]) else: raise ValueError(f"Frame type {frame_type} not Supported in PythonDEServer") diff --git a/deapi/tests/test_fake_server/test_gpu_frames.py b/deapi/tests/test_fake_server/test_gpu_frames.py index 618ba07d..21c10f5b 100644 --- a/deapi/tests/test_fake_server/test_gpu_frames.py +++ b/deapi/tests/test_fake_server/test_gpu_frames.py @@ -60,3 +60,13 @@ def test_stop_ends_the_simulated_frames(client): frames.close() assert 0 < seen < 8 * 8 * 1000 assert not client.acquiring + + +def test_external_image_is_served(client): + """A HAADF-only search reads external_image1; the simulated server answers it.""" + _scan(client, 8, 1) + client.start_acquisition(1) + result = client.get_result("external_image1", "FLOAT32") + image = getattr(result, "image", None) + image = result[0] if image is None else image + assert np.asarray(image).size > 0 From 595535c984b9000d3641c5a25bb278baea9d91db Mon Sep 17 00:00:00 2001 From: Carter Francis Date: Wed, 23 Sep 2026 14:50:10 -0500 Subject: [PATCH 3/4] Simulated server: wrap the scan position on repeated scans With Scan - Repeats an acquisition outlasts one pass, and the current navigation index (elapsed time x frame rate) ran past the scan: get_result raised mid-request and the client waited for a reply that never came. The index now wraps each pass and uses the same clamped frame rate as the acquisition's length. Claude-Session: https://claude.ai/code/session_017DXgzNPmF7pJt9JMGnKTrT --- deapi/simulated_server/fake_server.py | 7 ++++--- deapi/tests/test_fake_server/test_gpu_frames.py | 7 +++++-- 2 files changed, 9 insertions(+), 5 deletions(-) diff --git a/deapi/simulated_server/fake_server.py b/deapi/simulated_server/fake_server.py index 0f708170..90bc8344 100644 --- a/deapi/simulated_server/fake_server.py +++ b/deapi/simulated_server/fake_server.py @@ -269,10 +269,11 @@ def current_navigation_index(self): if self.acquisition_status == "Idle": return tuple(np.array(self.fake_data.navigator.shape) - 1) else: + # Frames so far at the rate the acquisition runs at, wrapped to the + # scan: with Scan - Repeats each pass starts again at the first position. + fps = min(float(self["Frames Per Second"]), SIMULATED_MAX_FPS) index = np.unravel_index( - int( - (time.time() - self.start_time) * float(self["Frames Per Second"]), - ), + int((time.time() - self.start_time) * fps) % self.fake_data.navigator.size, self.fake_data.navigator.shape, ) # only works for raster scans return index diff --git a/deapi/tests/test_fake_server/test_gpu_frames.py b/deapi/tests/test_fake_server/test_gpu_frames.py index 21c10f5b..511a44e7 100644 --- a/deapi/tests/test_fake_server/test_gpu_frames.py +++ b/deapi/tests/test_fake_server/test_gpu_frames.py @@ -63,10 +63,13 @@ def test_stop_ends_the_simulated_frames(client): def test_external_image_is_served(client): - """A HAADF-only search reads external_image1; the simulated server answers it.""" - _scan(client, 8, 1) + """A HAADF-only search reads external_image1; the simulated server answers it, + including after the first pass of a repeated scan.""" + _scan(client, 8, 1000) client.start_acquisition(1) + time.sleep(0.2) # past the first pass: 64 positions at 1000 fps result = client.get_result("external_image1", "FLOAT32") image = getattr(result, "image", None) image = result[0] if image is None else image assert np.asarray(image).size > 0 + client.stop_acquisition() From 3e4fc3f7912caa0e24e52cf0366eaa8a2bc3dd4c Mon Sep 17 00:00:00 2001 From: Carter Francis Date: Thu, 24 Sep 2026 17:46:40 -0500 Subject: [PATCH 4/4] Depend on anyplotlib from PyPI The dependency was anyplotlib from git (@main) plus a uv source pointing at a local ../anyplotlib checkout, so what got installed depended on whoever's working copy sat beside deapi, and a project depending on deapi inherited that. It is now anyplotlib>=0.10 from PyPI (0.10.1 today). deapi's fake-server, client and anyplotlib tests pass against 0.10.1, except test_live_result.py::test_idle_to_active_transition, which fails the same way on anyplotlib (the base branch) and with anyplotlib 0.7.3. Claude-Session: https://claude.ai/code/session_017DXgzNPmF7pJt9JMGnKTrT --- pyproject.toml | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index f09dd38e..32072ec8 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -24,7 +24,7 @@ dependencies = [ "scikit-image", "sympy", "pywin32; platform_system == 'Windows'", - "anyplotlib@git+https://github.com/cssfrancis/anyplotlib.git@main", + "anyplotlib>=0.10", ] description = "API for DE Server" version = "5.3b6" @@ -44,9 +44,6 @@ requires-python = ">=3.10" include = ["deapi*", "deapi.*"] where = ["."] -[tool.uv.sources] -anyplotlib = { path = "../anyplotlib", editable = true } - [project.optional-dependencies] tests = [ "pytest-instafail",