Merge branch 'main' of gitea.pkm.physik.tu-darmstadt.de:IPKM/python3-damaris
This commit is contained in:
markusro committed 2026-07-14 23:08:12 +02:00
commit 8bd088d97f
15 files changed
+2700 -232

No files matched your search

+12
View File
@@ -11,3 +11,15 @@ debian/files
.junie
debian/debhelper-build-stamp
debian/.debhelper
# Sphinx documentation
doc/_build/
# Debian build artifacts
debian/*.debhelper.log
debian/*.debhelper
debian/*.substvars
debian/python3-damaris/
# Agent metadata
AGENTS.md
+1
View File
@@ -2,5 +2,6 @@ include src/damaris/gui/DAMARIS3.png
include src/damaris/gui/DAMARIS3.ico
include src/damaris/gui/damaris3.glade
include src/damaris/gui/python.xml
include scripts/damaris-cli
#recursive-include doc/reference-html *.html *.gif *.png *.css
#recursive-include doc/tutorial-html *.html *.gif *.png *.css *.tar.gz *.sh
+1
View File
@@ -2,3 +2,4 @@ doc/ usr/share/doc/python3-damaris/
debian/damaris3.desktop usr/share/applications/
src/damaris/gui/DAMARIS3.png usr/share/pixmaps/
src/damaris/gui/DAMARIS3.ico usr/share/pixmaps/
scripts/damaris-cli usr/bin/
+3
View File
@@ -21,6 +21,9 @@ dependencies = [
"pyxdg",
"PyGObject >= 3.14.0, <= 3.51.0",
"pyserial>=3.5",
"sphinx>=9.0.4",
"sphinxcontrib-bibtex>=2.7.0",
"sphinx-rtd-theme>=3.1.0",
]
authors = [
{name = "Achim Gädke"},
+52 -32
View File
@@ -1,9 +1,11 @@
#! /usr/bin/env python
#! /usr/bin/env python3
import time
import sys
import os
import os.path
import argparse
import logging
import tables
import damaris.data.DataPool as DataPool
import damaris.gui.ResultReader as ResultReader
@@ -12,7 +14,9 @@ import damaris.gui.BackendDriver as BackendDriver
import damaris.gui.ResultHandling as ResultHandling
import damaris.gui.ExperimentHandling as ExperimentHandling
def some_listener(event):
logger = logging.getLogger("damaris.cli")
def progressbar(event):
if event.subject=="__recentexperiment" or event.subject=="__recentresult":
r=event.origin.get("__recentresult",-1)+1
e=event.origin.get("__recentexperiment",-1)+1
@@ -20,17 +24,19 @@ def some_listener(event):
ratio=100.0*r/e
else:
ratio=100.0
print("\r%d/%d (%.0f%%)"%(r,e,ratio), end=' ')
sys.stderr.write("\r%d/%d (%.0f%%)"%(r,e,ratio))
sys.stderr.flush()
class ScriptInterface:
def __init__(self, exp_script=None, res_script=None, backend_executable=None, spool_dir="spool"):
def __init__(self, exp_script=None, res_script=None, backend_executable=None, spool_dir="spool", log_callback=None):
self.exp_script=exp_script
self.res_script=res_script
self.backend_executable=backend_executable
self.spool_dir=os.path.abspath(spool_dir)
self.exp_handling=self.res_handling=None
self._log_callback = log_callback
self.exp_writer=self.res_reader=self.back_driver=None
if self.backend_executable is not None:
@@ -48,9 +54,9 @@ class ScriptInterface:
def runScripts(self):
# get script engines
if self.exp_script and self.exp_writer:
self.exp_handling=ExperimentHandling.ExperimentHandling(self.exp_script, self.exp_writer, self.data)
self.exp_handling=ExperimentHandling.ExperimentHandling(self.exp_script, self.exp_writer, self.data, self._log_callback)
if self.res_script and self.res_reader:
self.res_handling=ResultHandling.ResultHandling(self.res_script, self.res_reader, self.data)
self.res_handling=ResultHandling.ResultHandling(self.res_script, self.res_reader, self.data, self._log_callback)
# start them
if self.exp_handling: self.exp_handling.start()
@@ -73,27 +79,29 @@ class ScriptInterface:
if not self.exp_handling.is_alive():
self.exp_handling.join()
if self.exp_handling.raised_exception:
print(": experiment script failed at line %d (function %s): %s"%(self.exp_handling.location[0],
self.exp_handling.location[1],
self.exp_handling.raised_exception))
logger.error(": experiment script failed at line %d (function %s): %s",
self.exp_handling.location[0],
self.exp_handling.location[1],
self.exp_handling.raised_exception)
else:
print(": experiment script finished")
logger.info(": experiment script finished")
self.exp_handling = None
if self.res_handling is not None:
if not self.res_handling.is_alive():
self.res_handling.join()
if self.res_handling.raised_exception:
print(": result script failed at line %d (function %s): %s"%(self.res_handling.location[0],
self.res_handling.location[1],
self.res_handling.raised_exception))
logger.error(": result script failed at line %d (function %s): %s",
self.res_handling.location[0],
self.res_handling.location[1],
self.res_handling.raised_exception)
else:
print(": result script finished")
logger.info(": result script finished")
self.res_handling = None
if self.back_driver is not None:
if not self.back_driver.is_alive():
print(": backend finished")
logger.info(": backend finished")
self.back_driver=None
except KeyboardInterrupt:
@@ -120,33 +128,45 @@ class ScriptInterface:
dump_file=None
# todo
except Exception as e:
print("dump failed", e)
logger.error("Dump failed: %s", e)
if __name__=="__main__":
parser = argparse.ArgumentParser(
description="Run DAMARIS experiment and result scripts."
)
parser.add_argument("experiment_script", help="Path to the experiment script file")
parser.add_argument("result_script", help="Path to the result script file")
parser.add_argument(
"spool_dir", nargs="?", default=None,
help="Spool directory (defaults to current working directory)"
)
parser.add_argument(
"--backend", default="/usr/lib/damaris/backends/Mobilecore",
help="Path to the backend executable"
)
parser.add_argument(
"--output", default="data_pool.h5",
help="Output HDF5 file path (default: data_pool.h5)"
)
args = parser.parse_args()
if len(sys.argv)==1:
print("%s: data_handling_script [spool directory]"%sys.argv[0])
sys.exit(1)
if len(sys.argv)==3:
spool_dir=os.getcwd()
else:
spool_dir=sys.argv[3]
expscriptfile=open(sys.argv[1])
expscript=expscriptfile.read()
resscriptfile=open(sys.argv[2])
resscript=resscriptfile.read()
spool_dir = args.spool_dir if args.spool_dir else os.getcwd()
si=ScriptInterface(expscript, resscript,"/usr/lib/damaris/backends/Mobilecore", spool_dir)
with open(args.experiment_script) as f:
expscript = f.read()
with open(args.result_script) as f:
resscript = f.read()
si.data.register_listener(some_listener)
si = ScriptInterface(expscript, resscript, args.backend, spool_dir)
si.data.register_listener(progressbar)
si.runScripts()
si.waitForScriptsEnding()
si.dump_data("data_pool.h5")
si.dump_data(args.output)
si=None
si = None
+134
View File
@@ -0,0 +1,134 @@
// -*- C++ -*-
/**
* @file logging_helper.h
* @brief Lightweight logging helper for the DAMARIS C++ backend.
*
* Writes structured log messages to stdout in a format parseable by the
* Python FileTailerHandler:
*
* [LEVEL] message
*
* Levels: DEBUG, INFO, WARNING, ERROR, CRITICAL
*
* Usage in C++ code:
* #include "logging_helper.h"
* log_info("Starting backend, spool=%s", spool_dir.c_str());
* log_error("Failed to open device: %s", strerror(errno));
*/
#ifndef DAMARIS_LOGGING_HELPER_H
#define DAMARIS_LOGGING_HELPER_H
#include <cstdio>
#include <cstdarg>
#include <cstring>
#include <ctime>
#include <string>
namespace damaris {
namespace logging {
// ---------------------------------------------------------------------------
// Internal helpers
// ---------------------------------------------------------------------------
inline std::string timestamp() {
char buf[32];
time_t now = time(nullptr);
struct tm tm_buf;
#ifdef _WIN32
localtime_s(&tm_buf, &now);
#else
localtime_r(&now, &tm_buf);
#endif
strftime(buf, sizeof(buf), "%Y-%m-%d %H:%M:%S", &tm_buf);
return std::string(buf);
}
inline void log_message(int level, const char* fmt, ...) {
const char* level_str = "INFO";
switch (level) {
case 0: level_str = "DEBUG"; break;
case 1: level_str = "INFO"; break;
case 2: level_str = "WARNING"; break;
case 3: level_str = "ERROR"; break;
case 4: level_str = "CRITICAL"; break;
}
char msg[4096];
va_list args;
va_start(args, fmt);
vsnprintf(msg, sizeof(msg), fmt, args);
va_end(args);
fprintf(stdout, "[%s] %s\n", level_str, msg);
fflush(stdout);
}
// ---------------------------------------------------------------------------
// Convenience macros / inline functions
// ---------------------------------------------------------------------------
inline void log_debug(const char* fmt, ...) {
char msg[4096];
va_list args;
va_start(args, fmt);
vsnprintf(msg, sizeof(msg), fmt, args);
va_end(args);
fprintf(stdout, "[DEBUG] %s\n", msg);
fflush(stdout);
}
inline void log_info(const char* fmt, ...) {
char msg[4096];
va_list args;
va_start(args, fmt);
vsnprintf(msg, sizeof(msg), fmt, args);
va_end(args);
fprintf(stdout, "[INFO] %s\n", msg);
fflush(stdout);
}
inline void log_warning(const char* fmt, ...) {
char msg[4096];
va_list args;
va_start(args, fmt);
vsnprintf(msg, sizeof(msg), fmt, args);
va_end(args);
fprintf(stdout, "[WARNING] %s\n", msg);
fflush(stdout);
}
inline void log_error(const char* fmt, ...) {
char msg[4096];
va_list args;
va_start(args, fmt);
vsnprintf(msg, sizeof(msg), fmt, args);
va_end(args);
fprintf(stdout, "[ERROR] %s\n", msg);
fflush(stdout);
}
inline void log_critical(const char* fmt, ...) {
char msg[4096];
va_list args;
va_start(args, fmt);
vsnprintf(msg, sizeof(msg), fmt, args);
va_end(args);
fprintf(stdout, "[CRITICAL] %s\n", msg);
fflush(stdout);
}
} // namespace logging
} // namespace damaris
// ---------------------------------------------------------------------------
// Optional C-style macros for convenience
// ---------------------------------------------------------------------------
#define DAMARIS_LOG_DEBUG(fmt, ...) damaris::logging::log_debug(fmt, ##__VA_ARGS__)
#define DAMARIS_LOG_INFO(fmt, ...) damaris::logging::log_info(fmt, ##__VA_ARGS__)
#define DAMARIS_LOG_WARNING(fmt, ...) damaris::logging::log_warning(fmt, ##__VA_ARGS__)
#define DAMARIS_LOG_ERROR(fmt, ...) damaris::logging::log_error(fmt, ##__VA_ARGS__)
#define DAMARIS_LOG_CRITICAL(fmt, ...) damaris::logging::log_critical(fmt, ##__VA_ARGS__)
#endif // DAMARIS_LOGGING_HELPER_H
+28 -39
View File
@@ -5,12 +5,16 @@ import sys
import time
import re
import glob
import logging
from . import ExperimentWriter
from . import ResultReader
from .logging_utils import FileTailerHandler
import threading
import types
import signal
logger = logging.getLogger("damaris.backend_driver")
if sys.platform=="win32":
import winreg
@@ -40,7 +44,7 @@ class BackendDriver(threading.Thread):
try:
os.makedirs(os.path.abspath(self.spool_dir))
except OSError as e:
print(e)
logger.error("Could not create backend's spool directory %s: %s", self.spool_dir, e)
raise AssertionError("could not create backend's spool directory %s "%self.spool_dir)
# remove stale state filenames
@@ -62,10 +66,10 @@ class BackendDriver(threading.Thread):
if os.path.isdir("/proc/%d"%core_pid):
raise AssertionError("found backend with pid %d (state file %s) in same spool dir"%(core_pid,statefilename))
else:
print("removing stale backend state file", statefilename)
logger.warning("Removing stale backend state file: %s", statefilename)
os.remove(statefilename)
else:
print("todo: take care of existing backend state files")
logger.warning("TODO: take care of existing backend state files (Windows)")
self.result_reader = ResultReader.BlockingResultReader(self.spool_dir,
no=0,
@@ -81,21 +85,9 @@ class BackendDriver(threading.Thread):
self.raised_exception=None
def run(self):
# take care of older logfiles
self.core_output_filename=os.path.join(self.spool_dir,"logdata")
if os.path.isfile(self.core_output_filename):
i=0
max_logs=100
while os.path.isfile(self.core_output_filename+".%02d"%i):
i+=1
while (i>=max_logs):
i-=1
os.remove(self.core_output_filename+".%02d"%i)
for j in range(i):
os.rename(self.core_output_filename+".%02d"%(i-j-1),self.core_output_filename+".%02d"%(i-j))
os.rename(self.core_output_filename, self.core_output_filename+".%02d"%0)
# create logfile
self.core_output=open(self.core_output_filename,"w")
# Log file for backend output
self.core_output_filename = os.path.join(self.spool_dir, "logdata")
self.core_output = open(self.core_output_filename, "w")
# again look out for existing state files
state_files=glob.glob(os.path.join(self.spool_dir,"*.state"))
@@ -118,6 +110,10 @@ class BackendDriver(threading.Thread):
stdout=self.core_output,
stderr=self.core_output)
# Start tailing backend log file for unified logging
self._log_tailer = FileTailerHandler(self.core_output_filename)
self._log_tailer.start()
# wait till state file shows up
timeout=10
# to do: how should I know core's state name????!!!!!
@@ -126,19 +122,16 @@ class BackendDriver(threading.Thread):
while len(state_files)==0:
if timeout<0 or self.core_input is None or self.core_input.poll() is not None or self.quit_flag.isSet():
# look into core log file and include contents
print(timeout, self.core_input, self.quit_flag.isSet())
if self.core_input is not None:
print(self.core_input.poll())
logger.error("Backend failed to start: timeout=%s poll=%s", timeout, self.core_input.poll())
log_message=''
self.core_input=None
if os.path.isfile(self.core_output_filename):
# to do include log data
log_message='\n'+''.join(open(self.core_output_filename,"r", errors='replace').readlines()[:10])
if not log_message:
log_message=" no error message from core"
self.core_output.close()
self.raised_exception="no state file appeared or backend died away:"+log_message
print(self.raised_exception)
logger.error(self.raised_exception)
self.quit_flag.set()
return
time.sleep(0.05)
@@ -147,7 +140,7 @@ class BackendDriver(threading.Thread):
# save the one
if len(state_files)>1:
print("did find more than one state file, taking first one!")
logger.warning("Found more than one state file, taking first one!")
self.statefilename=state_files[0]
# read state file
@@ -182,28 +175,28 @@ class BackendDriver(threading.Thread):
backend_result=self.core_input.poll()
if backend_result is not None: break
if wait_loop_counter==10:
print("sending termination signal to backend process")
logger.warning("Sending termination signal to backend process")
self.send_signal("SIGTERM")
elif wait_loop_counter==20:
print("sending kill signal to backend process")
logger.warning("Sending kill signal to backend process")
self.send_signal("SIGKILL")
elif wait_loop_counter>30:
print("no longer waiting for backend shutdown")
logger.warning("No longer waiting for backend shutdown")
break
if backend_result is None:
print("backend dit not end properly, please stop it manually")
logger.error("Backend did not end properly, please stop it manually")
elif backend_result>0:
print("backend returned ", backend_result)
logger.error("Backend returned exit code %d", backend_result)
elif backend_result<0:
sig_name=[x for x in dir(signal) if x.startswith("SIG") and \
x[3]!="_" and \
(type(signal.__dict__[x])is int) and \
signal.__dict__[x]==-backend_result]
if sig_name:
print("backend was terminated by signal ",sig_name[0])
logger.error("Backend was terminated by signal %s", sig_name[0])
else:
print("backend was terminated by signal no",-backend_result)
logger.error("Backend was terminated by signal no %d", -backend_result)
self.core_input = None
self.core_pid = None
@@ -226,12 +219,6 @@ class BackendDriver(threading.Thread):
if os.path.isfile(resultfilename):
os.remove(resultfilename)
def get_messages(self):
# return pending messages
if self.core_output.tell()==os.path.getsize(self.core_output_filename):
return None
return self.core_output.read()
def restart_queue(self):
self.send_signal("SIGUSR1")
@@ -250,7 +237,7 @@ class BackendDriver(threading.Thread):
def send_signal(self, sig):
if self.core_pid is None:
print("BackendDriver.send_signal is called with core_pid=None")
logger.error("BackendDriver.send_signal called with core_pid=None")
return
try:
if sys.platform[:5]=="linux":
@@ -262,7 +249,7 @@ class BackendDriver(threading.Thread):
kill_command=os.path.join(cygwin_path,"bin","kill.exe")
os.popen("%s -%s %d"%(kill_command,sig,self.core_pid))
except OSError as e:
print("could not send signal %s to core: %s"%(sig, str(e)))
logger.error("Could not send signal %s to core: %s", sig, str(e))
def is_busy(self):
@@ -294,3 +281,5 @@ class BackendDriver(threading.Thread):
if self.core_output:
self.core_output.close()
self.core_output=None
if hasattr(self, '_log_tailer'):
self._log_tailer.stop()
+287 -155
View File
@@ -70,36 +70,9 @@ debug = False
# version info
from damaris import __version__
class logstream:
gui_log = None
text_log = sys.__stdout__
def write( self, message ):
# default for stdout and stderr
if self.gui_log is not None:
self.gui_log( message )
if debug or self.gui_log is None:
self.text_log.write( message )
self.text_log.flush( )
def __call__( self, message ):
self.write( message )
def flush(self):
pass
def __del__( self ):
pass
global log
log = logstream( )
sys.stdout = log
sys.stderr = log
ExperimentHandling.log = log
ResultHandling.log = log
ResultReader.log = log
ExperimentWriter.log = log
BackendDriver.log = log
DataPool.log = log
# Unified logging
import logging
logger = logging.getLogger("damaris.gui")
class LockFile:
@@ -109,17 +82,17 @@ class LockFile:
self.conn = sqlite3.connect(self.file)
self.cursor = self.conn.cursor()
self.id = None
print("Connected to db..")
logger.debug("Connected to db..")
def add_experiment(self):
print("Add experiment..")
logger.debug("Add experiment..")
self.id = str(uuid.uuid4())
self.cursor.execute("INSERT INTO damaris VALUES (?, ?)",(self.id, "started"))
self.conn.commit()
return self.id
def del_experiment(self, id):
print("Delete experiment..")
logger.debug("Delete experiment..")
self.cursor.execute('DELETE FROM damaris WHERE uuid=?', (id,))
self.conn.commit()
return True
@@ -166,6 +139,13 @@ class DamarisGUI:
self.log = LogWindow( self.xml_gui , self)
# Initialize unified logging with GUI output
from damaris.gui import setup_logging
self._logger = setup_logging(
level=logging.DEBUG if debug else logging.INFO,
gui_callback=self.log,
)
self.sw = ScriptWidgets( self.xml_gui, self )
self.toolbar_init( )
@@ -502,23 +482,23 @@ class DamarisGUI:
def start_experiment( self, widget, data=None ):
import time
overall_start = time.time()
print(f"DEBUG: start_experiment called at {overall_start}")
logger.debug("start_experiment called at %s", overall_start)
# something running?
if self.si is not None or self.state != DamarisGUI.Edit_State:
print("Experiment is still running or cleaning up!")
logger.error("Experiment is still running or cleaning up!")
return
# get config values:
config_start = time.time()
actual_config = self.config.get( )
print(f"DEBUG: Config load took {time.time() - config_start:.3f}s")
logger.debug("Config load took %.3fs", time.time() - config_start)
# get scripts and start script interface
gui_start = time.time()
self.sw.disable_editing( )
exp_script, res_script = self.sw.get_scripts( )
print(f"DEBUG: GUI operations took {time.time() - gui_start:.3f}s")
logger.debug("GUI operations took %.3fs", time.time() - gui_start)
if not actual_config[ "start_result_script" ]:
res_script = ""
if not actual_config[ "start_experiment_script" ]:
@@ -528,7 +508,7 @@ class DamarisGUI:
backend = ""
if (backend == "" and exp_script == "" and res_script == ""):
print("nothing to do...so doing nothing!")
logger.info("Nothing to do...so doing nothing!")
self.sw.enable_editing( )
return
@@ -546,11 +526,11 @@ class DamarisGUI:
ln = 0
if type( lo ) is not int:
lo = 0
print("Experiment Script: %s at line %d, col %d:" % (e.__class__.__name__, ln, lo))
logger.error("Experiment Script: %s at line %d, col %d", e.__class__.__name__, ln, lo)
if e.text != "":
print(e.text.rstrip())
print(" "*(e.offset-1)+"^")
print(e)
logger.error(e.text.rstrip())
logger.error(" "*(e.offset-1)+"^")
logger.error("%s", e)
res_code = None
if res_script != "":
@@ -564,14 +544,14 @@ class DamarisGUI:
ln = 0
if type( lo ) is not int:
lo = 0
print("Result script: %s at line %d, col %d:" % (e.__class__.__name__, ln, lo))
logger.error("Result script: %s at line %d, col %d", e.__class__.__name__, ln, lo)
if e.text != "":
print("\"%s\"" % e.text)
# print " "*(e.offset+1)+"^" # nice idea, but needs monospaced fonts
logger.error("\"%s\"", e.text)
# logger.error(" "*(e.offset+1)+"^" # nice idea, but needs monospaced fonts
pass
print(e)
logger.error("%s", e)
print(f"DEBUG: Script compilation took {time.time() - compile_start:.3f}s")
logger.debug("Script compilation took %.3fs", time.time() - compile_start)
# detect error
if (exp_script != "" and exp_code is None) or \
@@ -604,10 +584,10 @@ class DamarisGUI:
pass
# Use observe_data_pool(None) which is now optimized with fast TreeStore replacement
print(f"DEBUG: Data pool cleanup ({listener_count} listeners, {data_items} items)")
logger.debug("Data pool cleanup (%d listeners, %d items)", listener_count, data_items)
self.monitor.observe_data_pool( None )
print(f"DEBUG: Data cleanup took {time.time() - data_cleanup_start:.3f}s")
logger.debug("Data cleanup took %.3fs", time.time() - data_cleanup_start)
# set the text mark for hdf logging
if self.log.textbuffer.get_mark( "lastdumped" ) is None:
@@ -618,14 +598,12 @@ class DamarisGUI:
interval = 0.1
self.stop_experiment_flag.clear()
# Add timing debug info
import time
start_time = time.time()
print(f"DEBUG: Starting experiment at {start_time}")
logger.debug("Starting experiment at %s", start_time)
self.id = self.lockfile.add_experiment()
print(f"DEBUG: Lock file operation took {time.time() - start_time:.3f}s")
logger.debug("Lock file operation took %.3fs", time.time() - start_time)
# Add more timing points
lock_wait_start = time.time()
while not self.lockfile.am_i_next():
@@ -638,9 +616,9 @@ class DamarisGUI:
loop_run = 0
loop_run += interval
print(f"DEBUG: Lock wait took {time.time() - lock_wait_start:.3f}s")
logger.debug("Lock wait took %.3fs", time.time() - lock_wait_start)
if self.stop_experiment_flag.isSet():
if self.stop_experiment_flag.is_set():
#Experiment has been stopped, clean up and leave
self.experiment_script_statusbar_label.set_text("Experiment stopped")
self.state = DamarisGUI.Edit_State
@@ -662,9 +640,10 @@ class DamarisGUI:
res_code,
backend,
self.spool_dir,
log_callback=self.log,
clear_jobs=actual_config[ "del_jobs_after_execution" ],
clear_results=actual_config[ "del_results_after_processing" ] )
print(f"DEBUG: ScriptInterface creation took {time.time() - script_interface_start:.3f}s")
logger.debug("ScriptInterface creation took %.3fs", time.time() - script_interface_start)
self.data = self.si.data
# run frontend and script engines
@@ -672,18 +651,18 @@ class DamarisGUI:
# Start temperature monitoring if configured (before scripts start)
temp_start = time.time()
self.start_temperature_monitoring(actual_config)
print(f"DEBUG: Temperature monitoring startup took {time.time() - temp_start:.3f}s")
logger.debug("Temperature monitoring startup took %.3fs", time.time() - temp_start)
scripts_start = time.time()
self.si.runScripts( )
print(f"DEBUG: Script execution startup took {time.time() - scripts_start:.3f}s")
logger.debug("Script execution startup took %.3fs", time.time() - scripts_start)
except Exception as e:
#print "ToDo evaluate exception",str(e), "at",traceback.extract_tb(sys.exc_info()[2])[-1][1:3]
#print "Full traceback:"
#logger.debug("ToDo evaluate exception %s at %s", str(e), traceback.extract_tb(sys.exc_info()[2])[-1][1:3])
#logger.debug("Full traceback:")
traceback_file = io.StringIO( )
traceback.print_tb( sys.exc_info( )[ 2 ], None, traceback_file )
self.main_notebook.set_current_page( DamarisGUI.Log_Display )
print("Error while executing scripts: %s\n" % str( e ) + traceback_file.getvalue( ))
logger.error("Error while executing scripts: %s", str( e ) + traceback_file.getvalue( ))
traceback_file = None
self.data = None
@@ -691,14 +670,14 @@ class DamarisGUI:
still_running = [_f for _f in [ self.si.exp_handling, self.si.res_handling, self.si.back_driver ] if _f]
for r in still_running:
r.quit_flag.set( )
print("waiting for threads stoping...", end=' ')
logger.info("Waiting for threads to stop...")
still_running = [x for x in [ self.si.exp_handling, self.si.res_handling, self.si.back_driver ] if x is not None and x.is_alive( )]
for t in still_running:
while t.is_alive():
while gtk.events_pending():
gtk.main_iteration_do(False)
t.join( 0.05 )
print("done")
logger.info("Threads stopped")
# cleanup
self.si = None
@@ -785,17 +764,18 @@ class DamarisGUI:
if not self.si.exp_handling.is_alive( ):
self.si.exp_handling.join( )
if self.si.exp_handling.raised_exception:
print("experiment script failed at line %d (function %s): %s" % (self.si.exp_handling.location[ 0 ],
self.si.exp_handling.location[ 1 ],
self.si.exp_handling.raised_exception))
print("Full traceback", self.si.exp_handling.traceback)
logger.error("Experiment script failed at line %d (function %s): %s",
self.si.exp_handling.location[ 0 ],
self.si.exp_handling.location[ 1 ],
self.si.exp_handling.raised_exception)
logger.error("Full traceback: %s", self.si.exp_handling.traceback)
e_text = "Experiment Script Failed (%d)" % e
else:
e_text = "Experiment Script Finished (%d)" % e
print("experiment script finished")
logger.info("Experiment script finished")
self.si.exp_handling = None
else:
#print self.si.back_driver.get_messages()
# Backend messages are now handled by FileTailerHandler
e_text = "Experiment Script Running (%d)" % e
e_text += experimenttimetext
@@ -803,10 +783,11 @@ class DamarisGUI:
if not self.si.res_handling.is_alive( ):
self.si.res_handling.join( )
if self.si.res_handling.raised_exception:
print("result script failed at line %d (function %s): %s" % (self.si.res_handling.location[ 0 ],
self.si.res_handling.location[ 1 ],
self.si.res_handling.raised_exception))
print("Full traceback", self.si.res_handling.traceback)
logger.error("Result script failed at line %d (function %s): %s",
self.si.res_handling.location[ 0 ],
self.si.res_handling.location[ 1 ],
self.si.res_handling.raised_exception)
logger.error("Full traceback: %s", self.si.res_handling.traceback)
r_text = "Result Script Failed (%d)" % r
else:
r_text = "Result Script Finished (%d)" % r
@@ -830,16 +811,16 @@ class DamarisGUI:
if self.dump_thread is not None:
if self.dump_thread.is_alive( ):
sys.stdout.write( "." )
sys.stderr.write( "." )
self.dump_dots += 1
if self.dump_dots > 80:
print()
sys.stderr.write("\n")
self.dump_dots = 0
else:
self.dump_thread.join( )
self.dump_thread = None
dump_size = os.stat( self.dump_filename ).st_size / 1e6
print("done (%.1f s, %.1f MB)" % (time.time( ) - self.dump_start_time, dump_size))
logger.info("Done (%.1f s, %.1f MB)", time.time( ) - self.dump_start_time, dump_size)
gdk.threads_enter( )
if e_text:
@@ -855,12 +836,12 @@ class DamarisGUI:
# Stop temperature monitoring when experiment finishes
self.stop_temperature_monitoring()
if self.save_thread is None and self.dump_filename != "":
print("all subprocesses ended, saving data pool")
logger.info("All subprocesses ended, saving data pool")
# Disable run button during data saving to prevent UI confusion
gdk.threads_enter()
self.toolbar_run_button.set_sensitive(False)
gdk.threads_leave()
# thread to save data...
self.save_thread = threading.Thread( target=self.dump_states, name="dump states" )
self.save_thread.start( )
@@ -874,21 +855,21 @@ class DamarisGUI:
self.toolbar_stop_button.set_sensitive( False )
gdk.threads_leave( )
if len( still_running ) != 0:
print("subprocess(es) still running: " + ', '.join( [s.getName( ) for s in still_running] ))
logger.warning("Subprocess(es) still running: %s", ', '.join( [s.getName( ) for s in still_running] ))
return True
else:
if self.save_thread is not None:
if self.save_thread.is_alive( ):
sys.stdout.write( "." )
sys.stderr.write( "." )
self.dump_dots += 1
if self.dump_dots > 80:
print()
sys.stderr.write("\n")
self.dump_dots = 0
return True
self.save_thread.join( )
self.save_thread = None
dump_size = os.stat( self.dump_filename ).st_size / 1e6
print("done (%.1f s, %.1f MB)" % (time.time( ) - self.dump_start_time, dump_size))
logger.info("Done (%.1f s, %.1f MB)", time.time( ) - self.dump_start_time, dump_size)
# keep data to display but throw away everything else
self.si = None
@@ -909,7 +890,7 @@ class DamarisGUI:
dump_interval = self.dump_timeinterval if self.dump_timeinterval > 0 else 2.0
if self.dump_thread is None and self.dump_filename != "" and \
self.last_dumped + dump_interval < time.time( ):
print("dumping data pool")
logger.info("Dumping data pool")
self.dump_start_time = time.time( )
self.dump_dots = 0
# thread to save data...
@@ -944,13 +925,13 @@ class DamarisGUI:
try:
self.dump_filename = actual_config.get( "data_pool_name" ) % extensions_dict
except ValueError:
print("invalid formating character in '" + actual_config.get( "data_pool_name" ) + "'")
logger.error("Invalid formatting character in '%s'", actual_config.get( "data_pool_name" ))
self.dump_filename = "DAMARIS_data_pool.h5"
print("reseting dumpfile to " + self.dump_filename)
logger.info("Resetting dumpfile to %s", self.dump_filename)
del extensions_dict
if os.path.split( self.dump_filename )[ 1 ] == "" or os.path.isdir( self.dump_filename ):
print("The dump filename is a directory, using filename 'DAMARIS_data_pool.h5'")
logger.warning("Dump filename is a directory, using filename 'DAMARIS_data_pool.h5'")
self.dump_filename += os.sep + "DAMARIS_data_pool.h5"
# compression
@@ -967,7 +948,7 @@ class DamarisGUI:
try:
self.dump_timeinterval = int( 60 * float( actual_config[ "data_pool_write_interval" ] ) )
except ValueError as e:
print("configuration provides non-number for dump interval: " + str( e ))
logger.warning("Configuration provides non-number for dump interval: %s", str( e ))
# if existent, move away old files
if os.path.isfile( self.dump_filename ):
@@ -989,15 +970,15 @@ class DamarisGUI:
last_backup -= 1
os.rename( self.dump_filename, dump_filename_pattern % 0 )
if cummulated_size > (1 << 30):
print("Warning: the cumulated backups size of '%s' is %d MByte" % (self.dump_filename,
cummulated_size / (1 << 20)))
logger.warning("The cumulated backups size of '%s' is %d MByte", self.dump_filename,
cummulated_size / (1 << 20))
# init is finnished now
# now it's time to create the hdf file
dump_file = None
if not os.path.isfile( self.dump_filename ):
if not init:
print("dump file \"%s\" vanished unexpectedly, creating new one" % self.dump_filename)
logger.warning("Dump file \"%s\" vanished unexpectedly, creating new one", self.dump_filename)
# have a look to the path and create necessary directories
dir_stack = [ ]
dir_trunk = os.path.dirname( os.path.abspath( self.dump_filename ) )
@@ -1011,8 +992,7 @@ class DamarisGUI:
continue
os.mkdir( dir_trunk )
except OSError as e:
print(e)
print("could not create dump directory '%s', so hdf5 dumps disabled" % dir_trunk)
logger.error("Could not create dump directory '%s', so hdf5 dumps disabled: %s", dir_trunk, e)
self.dump_filename = ""
self.dump_timeinterval = 0
return True
@@ -1055,7 +1035,7 @@ class DamarisGUI:
if dump_file is None:
# exit!
print("could not create dump directory '%s', so hdf5 dumps disabled" % dir_trunk)
logger.error("Could not create dump directory '%s', so hdf5 dumps disabled", dir_trunk)
self.dump_filename = ""
self.dump_timeinterval = 0
return True
@@ -1199,24 +1179,24 @@ class DamarisGUI:
temp_baudrate = config.get("temperature_baudrate", 19200)
if not temp_device:
print("Temperature monitoring: No device configured for Eurotherm")
logger.warning("Temperature monitoring: No device configured for Eurotherm")
return
self.temp_controller = create_controller(
"eurotherm", temp_device, baudrate=temp_baudrate
)
print(f"Temperature monitoring started with Eurotherm on {temp_device} (interval: {temp_interval}s)")
logger.info("Temperature monitoring started with Eurotherm on %s (interval: %ss)", temp_device, temp_interval)
elif controller_type == "simulated":
# Simulated controller doesn't need device configuration
# Use default parameters: 290K base, 10K amplitude, 60s period, 5% noise, 10% missing
self.temp_controller = create_controller("simulated")
print(f"Temperature monitoring started with simulated controller (interval: {temp_interval}s)")
logger.info("Temperature monitoring started with simulated controller (interval: %ss)", temp_interval)
else:
print(f"Temperature monitoring: Unknown controller type: {controller_type}")
logger.warning("Temperature monitoring: Unknown controller type: %s", controller_type)
return
# Create and start monitor thread
self.temp_monitor = TemperatureMonitor(
controller=self.temp_controller,
@@ -1225,9 +1205,9 @@ class DamarisGUI:
data_key="temperature"
)
self.temp_monitor.start()
except Exception as e:
print(f"Failed to start temperature monitoring: {e}")
logger.error("Failed to start temperature monitoring: %s", e)
import traceback
traceback.print_exc()
self.temp_controller = None
@@ -1240,17 +1220,17 @@ class DamarisGUI:
if self.temp_monitor is not None:
try:
self.temp_monitor.stop()
print("Temperature monitoring stopped")
logger.info("Temperature monitoring stopped")
except Exception as e:
print(f"Error stopping temperature monitoring: {e}")
logger.error("Error stopping temperature monitoring: %s", e)
finally:
self.temp_monitor = None
if self.temp_controller is not None:
try:
self.temp_controller.close()
except Exception as e:
print(f"Error closing temperature controller: {e}")
logger.error("Error closing temperature controller: %s", e)
finally:
self.temp_controller = None
@@ -1281,7 +1261,7 @@ class DamarisGUI:
if not self.doc_browser.is_alive( ):
self.doc_browser.join( )
if self.doc_browser.my_webbrowser is not None:
print("new browser tab")
logger.debug("New browser tab")
self.doc_browser.my_webbrowser.open_new_tab( self.doc_urls[ requested_doc ] )
else:
del self.doc_browser
@@ -1291,7 +1271,7 @@ class DamarisGUI:
self.doc_browser = start_browser( self.doc_urls[ requested_doc ] )
self.doc_browser.start( )
else:
print("missing docs for '%s'" % (requested_doc))
logger.warning("Missing docs for '%s'", requested_doc)
show_manual = show_doc_menu
@@ -1315,11 +1295,11 @@ class start_browser( threading.Thread ):
# this is what it should be everywhere!
self.my_webbrowser = webbrowser.get( )
if self.my_webbrowser is not None:
print("starting web browser (module webbrowser)")
logger.info("Starting web browser (module webbrowser)")
self.my_webbrowser.open( self.start_url )
return True
# last resort
print("starting web browser (webbrowser.py)")
logger.info("Starting web browser (webbrowser.py)")
self.my_webbrowser_process = os.spawnl( os.P_NOWAIT,
sys.executable,
os.path.basename( sys.executable ),
@@ -1329,8 +1309,23 @@ class start_browser( threading.Thread ):
class LogWindow:
"""
writes messages to the log window
"""Manages the application log window.
Provides a timestamped, scrollable text view for application messages.
Messages are dispatched via ``gobject.idle_add`` to keep the GTK main
loop responsive.
Args:
xml_gui: The :class:`gtk.Builder` instance containing the GUI widgets.
damaris_gui: Reference to the main :class:`DamarisGUI` application.
Attributes:
textview (:class:`gtk.TextView`): The monospaced text view widget.
textbuffer (:class:`gtk.TextBuffer`): The buffer backing the text view.
Examples:
>>> log = LogWindow(xml_gui, damaris_gui)
>>> log("Application started")
"""
def __init__( self, xml_gui, damaris_gui ):
@@ -1342,8 +1337,6 @@ class LogWindow:
font_desc = pango.FontDescription("Monospace 10")
self.textview.modify_font(font_desc)
self.textbuffer = self.textview.get_buffer( )
self.logstream = log
self.logstream.gui_log = self
self.last_timetag = None
self( "Started in directory %s\n" % os.getcwd( ) )
@@ -1373,14 +1366,59 @@ class LogWindow:
return False
def __del__( self ):
self.logstream.gui_log = None
self.logstream = None
pass
class ScriptWidgets:
"""Manages experiment and result script editing widgets.
Provides two syntax-highlighted, auto-completing text views for editing
DAMARIS experiment and result scripts. Handles file open/save dialogs,
toolbar state management, line/column indicators, and search functionality.
Args:
xml_gui: The :class:`gtk.Builder` instance containing the GUI widgets.
outer_space: Reference to the parent application (typically
:class:`DamarisGUI`) that owns this widget set.
Attributes:
editing_state (bool): ``True`` if scripts are editable, ``False`` during
experiment execution.
exp_script_filename (str or None): Path to the currently loaded
experiment script file.
res_script_filename (str or None): Path to the currently loaded
result script file.
experiment_script_textview (:class:`gtksourceview2.View`): Text view
for the experiment script.
data_handling_textview (:class:`gtksourceview2.View`): Text view for
the result script.
Public Methods:
set_scripts: Replace script buffer contents.
get_scripts: Retrieve current script text from both buffers.
enable_editing: Re-enable script editing.
disable_editing: Lock scripts during experiment execution.
open_file: Open a script file via file chooser dialog.
open_hdf5: Load data and scripts from an HDF5 file.
save_file: Save the current script to its associated file.
save_file_as: Save the current script with a new filename.
save_all_files: Save both experiment and result scripts.
new_file: Clear the current script buffer.
check_script: Compile-check the current script for syntax errors.
undo / redo: Undo/redo text edits.
search: Open the search-and-replace dialog.
"""
def __init__( self, xml_gui, outer_space ):
"""
initialize text widgets with text
"""Initialize text widgets with syntax highlighting and completion.
Sets up two :class:`gtksourceview2.View` instances with Python
syntax highlighting, keyword auto-completion, line numbers, and
smart indentation. Connects text buffer signals to toolbar and
status bar update handlers.
Args:
xml_gui: The :class:`gtk.Builder` instance containing the GUI widgets.
outer_space: Reference to the parent application.
"""
self.xml_gui = xml_gui
self.outer_space = outer_space
@@ -1641,8 +1679,8 @@ get_ylabel set_ylabel"""
ln = 0
if type( lo ) is not int:
lo = 0
print("Syntax Error:\n%s in %s at line %d, offset %d" % (
str( se ), se.filename, ln, lo) + "\n(ToDo: Dialog)")
logger.error("Syntax Error:\n%s in %s at line %d, offset %d",
str( se ), se.filename, ln, lo)
if ln > 0 and ln <= tb.get_line_count( ):
new_place = tb.get_iter_at_line_offset( ln - 1, 0 )
if lo > 0 and lo <= new_place.get_chars_in_line( ):
@@ -1650,7 +1688,7 @@ get_ylabel set_ylabel"""
tb.place_cursor( new_place )
tv.scroll_to_iter( new_place, 0.2, False, 0, 0 )
except Exception as e:
print("Compilation Error:\n" + str( e ) + "\n(ToDo: Dialog)")
logger.error("Compilation Error:\n%s", str( e ))
def notebook_page_switched( self, notebook, page, pagenumber ):
self.set_toolbuttons_status( )
@@ -1668,7 +1706,7 @@ get_ylabel set_ylabel"""
newline = self.data_handling_line_indicator.get_value_as_int( ) - 1
newcol = column_indicator.get_value_as_int( ) - 1
else:
print("unknown line/column selector")
logger.error("Unknown line/column selector")
return False
textbuffer = textview.get_buffer( )
@@ -2474,13 +2512,13 @@ pygobject version %(pygobject)s
elif isinstance( model, gtk.TreeStore ):
model.append( None, [ libname ] )
else:
print("cannot append compression lib name to %s" % model.__class__.__name__)
logger.warning("Cannot append compression lib name to %s", model.__class__.__name__)
# debug message
if debug:
print("DAMARIS", __version__)
print(components_text % components_versions)
logger.debug("DAMARIS %s", __version__)
logger.debug(components_text % components_versions)
# set no compression as default...
self.config_data_pool_complib.set_active( 0 )
@@ -2496,7 +2534,7 @@ pygobject version %(pygobject)s
if os.access( self.system_default_filename, os.R_OK ):
self.load_config( self.system_default_filename )
else:
print("can not read system defaults from %s, ask your instrument responsible if required" % self.system_default_filename)
logger.warning("Can not read system defaults from %s, ask your instrument responsible if required", self.system_default_filename)
self.config_from_system = self.get( )
self.load_config( )
@@ -2619,7 +2657,7 @@ pygobject version %(pygobject)s
iter = model.iter_next( iter )
# if this compression method is not supported, warn and do nothing
if iter is None:
print("compression method %s is not supported" % config[ "data_pool_complib" ])
logger.warning("Compression method %s is not supported", config[ "data_pool_complib" ])
def on_adc_bit_depth_changed(self, widget):
ADC_Result.default_bit_depth = widget.get_value_as_int()
@@ -2783,7 +2821,7 @@ pygobject version %(pygobject)s
readfile = open( filename, "rb" )
except Exception as e:
if debug:
print("Could not open %s: %s" % (filename, str( e )))
logger.debug("Could not open %s: %s", filename, str( e ))
return
# parser functions
@@ -2835,7 +2873,7 @@ pygobject version %(pygobject)s
if not os.path.isdir( dirs ):
os.makedirs( dirs )
print("save config to: "+filename)
logger.info("Save config to: %s", filename)
configfile = open( filename, "w" )
configfile.write( "<?xml version='1.0'?>\n" )
@@ -2844,10 +2882,10 @@ pygobject version %(pygobject)s
if k in self.config_from_system \
and self.config_from_system[ k ] == v:
if debug:
print("Ignoring for write, because system value for %r is %r equal to %r" % \
(k, self.config_from_system[ k ], v))
logger.debug("Ignoring for write, because system value for %r is %r equal to %r",
k, self.config_from_system[ k ], v)
continue
print(k,v, type(v))
logger.debug("Config key: %s = %r (%s)", k, v, type(v).__name__)
val = ""
typename = ""
if type( v ) is bool:
@@ -2870,9 +2908,54 @@ pygobject version %(pygobject)s
class MonitorWidgets:
"""Manages the matplotlib-based data visualization panel.
Provides an interactive plot canvas with zoom, pan, scaling controls,
a data source selector, and real-time updates from the DAMARIS data pool.
Supports :class:`Accumulation`, :class:`ADC_Result`, and
:class:`MeasurementResult` data types with automatic axis rescaling,
error bars, and statistics overlays.
Args:
xml_gui: The :class:`gtk.Builder` instance containing the GUI widgets.
Attributes:
matplot_figure (:class:`matplotlib.figure.Figure`): The underlying
matplotlib figure.
matplot_axes (:class:`matplotlib.axes.Axes`): The plot axes.
matplot_canvas (:class:`matplotlib.backends.backend_gtk3.FigureCanvas`):
The GTK canvas embedding the figure.
matplot_toolbar (:class:`matplotlib.backends.backend_gtk3.NavigationToolbar2GTK3`):
The navigation toolbar with pan/zoom/cursor controls.
data_pool (:class:`damaris.data.DataPool` or None): The data pool
being observed for real-time updates.
displayed_data (list): ``[data_name, data_object]`` of the currently
selected data, or ``[None, None]``.
Public Methods:
observe_data_pool: Register a listener and begin observing a data pool.
source_list_reset: Reset the data source selector tree.
source_list_add: Add a data source entry to the selector tree.
source_list_remove: Remove a data source entry from the selector tree.
source_list_current: Get the currently selected data source name.
clear_display: Clear all plot elements and reset axes.
update_display: Redraw the plot with current data.
renew_display: Clear and redraw the plot with current data.
save_display_data_as_text: Export displayed data to a CSV file.
toggle_cursor: Toggle the matplotlib crosshair cursor.
copy_canvas_as_png: Copy the plot as a PNG image to the clipboard.
copy_canvas_as_svg: Copy the plot as SVG text to the clipboard.
"""
def __init__( self, xml_gui ):
"""
initialize matplotlib widgets and stuff around
"""Initialize matplotlib widgets, canvas, toolbar, and source selector.
Creates a :class:`matplotlib.figure.Figure` with linear x/y axes,
embeds it in a GTK canvas, attaches a navigation toolbar with a
custom crosshair toggle button, and builds a hierarchical source
selector :class:`gtk.TreeStore`.
Args:
xml_gui: The :class:`gtk.Builder` instance containing the GUI widgets.
"""
self.xml_gui = xml_gui
@@ -3096,11 +3179,11 @@ class MonitorWidgets:
pwd = namelist[ : ]
iter = self.source_list_find( namelist )
if iter is None or len( namelist ) > 0:
print("source_list_remove: WARNING: Not found")
logger.warning("source_list_remove: Not found")
return
model = self.display_source_treestore
if model.iter_has_child( iter ):
print("source_list_remove: WARNING: Request to delete a tree")
logger.warning("source_list_remove: Request to delete a tree")
return
while True:
parent = model.iter_parent( iter )
@@ -3185,7 +3268,7 @@ class MonitorWidgets:
if event.subject.startswith( "__" ):
return
if debug and self.update_counter < 0:
print("negative event count!", self.update_counter)
logger.debug("Negative event count: %d", self.update_counter)
# Throttling removed as it blocks data production thread.
# GUI updates are already decoupled via idle_add.
@@ -3231,10 +3314,10 @@ class MonitorWidgets:
do fast work selecting important events
"""
if debug and self.update_counter < 0:
print("negative event count!", self.update_counter)
logger.debug("Negative event count: %d", self.update_counter)
if self.update_counter > 5:
if debug:
print("sleeping to find time for graphics updates")
logger.debug("Sleeping to find time for graphics updates")
threading.Event( ).wait( 0.05 )
while self.update_counter > 15:
threading.Event( ).wait( 0.05 )
@@ -3263,7 +3346,7 @@ class MonitorWidgets:
if self.displayed_data[ 1 ] is new_data_struct:
# update display only
if self.update_counter > 10:
print("update queue too long (%d>10): skipping one update" % self.update_counter)
logger.warning("Update queue too long (%d>10): skipping one update", self.update_counter)
else:
gdk.threads_enter( )
try:
@@ -3281,7 +3364,7 @@ class MonitorWidgets:
new_data_struct.register_listener( self.datastructures_listener )
self.displayed_data[ 1 ] = new_data_struct
if self.update_counter > 10:
print("update queue too long (%d>10): skipping one update" % self.update_counter)
logger.warning("Update queue too long (%d>10): skipping one update", self.update_counter)
else:
gdk.threads_enter( )
try:
@@ -3322,7 +3405,7 @@ class MonitorWidgets:
self.update_counter_lock.release( )
# print "update display", self.update_counter
if self.update_counter > 10:
print("update queue too long (%d>10): skipping one update" % self.update_counter)
logger.warning("Update queue too long (%d>10): skipping one update", self.update_counter)
return
if self.displayed_data[ 0 ] is None or subject != self.displayed_data[ 0 ]:
return
@@ -3750,7 +3833,7 @@ class MonitorWidgets:
moving_average = data_slice = None
if max_points_to_display > 0 and len( xdata ) > max_points_to_display:
print("decimating data to %d points by moving average (prevent crash of matplotlib)" % max_points_to_display)
logger.warning("Decimating data to %d points by moving average (prevent crash of matplotlib)", max_points_to_display)
n = numpy.ceil( len( xdata ) / max_points_to_display )
moving_average = numpy.ones( n, dtype="float" ) / n
data_slice = numpy.array( numpy.floor( numpy.arange( max_points_to_display, dtype="float" ) \
@@ -3975,14 +4058,55 @@ class MonitorWidgets:
class ScriptInterface:
"""
texts or code objects are executed as experiment and result script the backend is started with sufficient arguments
"""Executes experiment and result scripts, coordinating with the DAMARIS backend.
Manages the lifecycle of experiment and result script execution. Depending
on whether a backend executable is provided, it either uses
:class:`BackendDriver.BackendDriver` to spawn an external process or falls
back to in-process :class:`ExperimentWriter` and :class:`ResultReader`
components. A shared :class:`DataPool` is used for communication between
scripts and the GUI.
Args:
exp_script (str or None): Python source code for the experiment script.
res_script (str or None): Python source code for the result script.
backend_executable (str or None): Path to the DAMARIS backend executable.
If ``None``, in-process execution is used.
spool_dir (str): Directory for inter-process communication spool files.
Defaults to ``"spool"``.
clear_jobs (bool): Clear pending jobs on startup. Defaults to ``True``.
clear_results (bool): Clear pending results on startup. Defaults to ``True``.
log_callback (callable or None): Optional callback for log messages.
Attributes:
data (:class:`damaris.data.DataPool`): Shared data pool for
experiment/result communication.
exp_handling (:class:`ExperimentHandling.ExperimentHandling` or None):
The running experiment handler thread.
res_handling (:class:`ResultHandling.ResultHandling` or None):
The running result handler thread.
back_driver (:class:`BackendDriver.BackendDriver` or None):
The backend process driver, if a backend executable was provided.
Public Methods:
runScripts: Start experiment, result, and backend handlers.
"""
def __init__( self, exp_script=None, res_script=None, backend_executable=None, spool_dir="spool", clear_jobs=True,
clear_results=True ):
"""
run experiment scripts and result scripts
clear_results=True, log_callback=None ):
"""Initialize the script execution environment.
Sets up the data pool, experiment writer, result reader, and optional
backend driver based on the provided configuration.
Args:
exp_script (str or None): Python source code for the experiment script.
res_script (str or None): Python source code for the result script.
backend_executable (str or None): Path to the DAMARIS backend executable.
spool_dir (str): Directory for inter-process communication spool files.
clear_jobs (bool): Clear pending jobs on startup.
clear_results (bool): Clear pending results on startup.
log_callback (callable or None): Optional callback for log messages.
"""
self.exp_script = exp_script
@@ -3992,6 +4116,7 @@ class ScriptInterface:
self.clear_jobs = clear_jobs
self.clear_results = clear_results
self.exp_handling = self.res_handling = None
self._log_callback = log_callback
self.exp_writer = self.res_reader = self.back_driver = None
if self.backend_executable is not None and self.backend_executable != "":
@@ -4016,13 +4141,20 @@ class ScriptInterface:
self.data = DataPool( )
def runScripts( self ):
"""Start experiment, result, and backend handler threads.
Creates and starts :class:`ExperimentHandling` and
:class:`ResultHandling` threads, and optionally starts the backend
driver process. Waits for the backend to report a valid PID before
proceeding.
"""
try:
# get script engines
self.exp_handling = self.res_handling = None
if self.exp_script and self.exp_writer:
self.exp_handling = ExperimentHandling.ExperimentHandling( self.exp_script, self.exp_writer, self.data )
self.exp_handling = ExperimentHandling.ExperimentHandling( self.exp_script, self.exp_writer, self.data, self._log_callback )
if self.res_script and self.res_reader:
self.res_handling = ResultHandling.ResultHandling( self.res_script, self.res_reader, self.data )
self.res_handling = ResultHandling.ResultHandling( self.res_script, self.res_reader, self.data, self._log_callback )
# start them
if self.back_driver is not None:
@@ -4031,7 +4163,7 @@ class ScriptInterface:
(self.back_driver.core_pid is None or self.back_driver.core_pid <= 0)):
while gtk.events_pending():
gtk.main_iteration_do(False)
self.back_driver.quit_flag.wait( 0.1 )
self.back_driver.quit_flag.wait( 0.05 )
if self.exp_handling:
self.exp_handling.start( )
if self.res_handling:
+40 -2
View File
@@ -3,24 +3,51 @@ import io
import traceback
import sys
import time
import logging
import ast
from damaris.experiments.Experiment import Quit
from damaris.experiments import Experiment
logger = logging.getLogger("damaris.experiment_handling")
class StopExperiment(Exception):
pass
class _SandboxStdout:
"""Redirects print() from sandboxed scripts to the GUI textview."""
def __init__(self, callback):
self._callback = callback
self._buffer = ""
def write(self, message):
self._buffer += message
if "\n" in self._buffer:
lines = self._buffer.split("\n")
self._buffer = lines[-1]
for line in lines[:-1]:
self._callback(line + "\n")
def flush(self):
if self._buffer:
self._callback(self._buffer)
self._buffer = ""
def isatty(self):
return False
class ExperimentHandling(threading.Thread):
"""
runs the experiment script in sandbox
"""
def __init__(self, script, exp_writer, data):
def __init__(self, script, exp_writer, data, log_callback=None):
threading.Thread.__init__(self, name="experiment handler")
self.script=script
self.writer=exp_writer
self.data=data
self.quit_flag = threading.Event()
self._log_callback = log_callback
if self.data is not None:
self.data["__recentexperiment"]=-1
@@ -50,6 +77,13 @@ class ExperimentHandling(threading.Thread):
self.location = None
exp_iterator=None
# Redirect stdout/stderr so print() in sandboxed scripts goes to the GUI
_orig_stdout = sys.stdout
_orig_stderr = sys.stderr
if self._log_callback is not None:
sys.stdout = _SandboxStdout(self._log_callback)
sys.stderr = _SandboxStdout(self._log_callback)
# check for time.sleep()
try:
tree = ast.parse(self.script)
@@ -66,7 +100,7 @@ class ExperimentHandling(threading.Thread):
found_blocking_sleep = True
if found_blocking_sleep:
print("Warning: time.sleep() detected in experiment script. This is blocking and not interruptible. Please use sleep() instead.")
logger.warning("time.sleep() detected in experiment script. This is blocking and not interruptible. Please use sleep() instead.")
break
except Exception:
pass
@@ -147,3 +181,7 @@ class ExperimentHandling(threading.Thread):
dataspace=None
self.exp_iterator=None
self.writer=None
# Restore stdout/stderr
sys.stdout = _orig_stdout
sys.stderr = _orig_stderr
+18 -3
View File
@@ -4,25 +4,36 @@ import sys
import os
import os.path
import traceback
import logging
import ast
from damaris.data import Resultable
from damaris.gui.ExperimentHandling import StopExperiment
from damaris.gui.ExperimentHandling import StopExperiment, _SandboxStdout
logger = logging.getLogger("damaris.result_handling")
class ResultHandling(threading.Thread):
"""
runs the result script in sandbox
"""
def __init__(self, script_data, result_iterator, data_pool):
def __init__(self, script_data, result_iterator, data_pool, log_callback=None):
threading.Thread.__init__(self,name="result handler")
self.script=script_data
self.results=result_iterator
self.data_space=data_pool
self.quit_flag=self.results.quit_flag
self._log_callback = log_callback
if self.data_space is not None:
self.data_space["__recentresult"]=-1
def run(self):
# Redirect stdout/stderr so print() in sandboxed scripts goes to the GUI
_orig_stdout = sys.stdout
_orig_stderr = sys.stderr
if self._log_callback is not None:
sys.stdout = _SandboxStdout(self._log_callback)
sys.stderr = _SandboxStdout(self._log_callback)
# execute it
dataspace={}
data_classes = __import__('damaris.data', dataspace, dataspace, ['*'])
@@ -52,7 +63,7 @@ class ResultHandling(threading.Thread):
found_blocking_sleep = True
if found_blocking_sleep:
print("Warning: time.sleep() detected in result script. This is blocking and not interruptible. Please use sleep() instead.")
logger.warning("time.sleep() detected in result script. This is blocking and not interruptible. Please use sleep() instead.")
break
except Exception:
pass
@@ -85,6 +96,10 @@ class ResultHandling(threading.Thread):
traceback_file=None
dataspace=None
# Restore stdout/stderr
sys.stdout = _orig_stdout
sys.stderr = _orig_stderr
def sleep(self, seconds):
self.quit_flag.wait(seconds)
if self.quit_flag.isSet():
+1 -1
View File
@@ -534,7 +534,7 @@ class BlockingResultReader(ResultReader):
def __init__(self, spool_dir=".", no=0, result_pattern="job.%09d.result", clear_jobs=False, clear_results=False):
ResultReader.__init__(self, spool_dir, no, result_pattern, clear_jobs=clear_jobs, clear_results=clear_results)
self.stop_no=None # end of job queue
self.poll_time=0.1 # sleep interval for polling results, <0 means no polling and stop
self.poll_time=0.05 # sleep interval for polling results, <0 means no polling and stop
self.in_advance=0
def __iter__(self):
+5
View File
@@ -0,0 +1,5 @@
from .logging_utils import (
FileTailerHandler,
LogDispatcher,
setup_logging,
)
+256
View File
@@ -0,0 +1,256 @@
# -*- coding: utf-8 -*-
"""
Unified logging for DAMARIS.
Provides:
- FileTailerHandler: reads new lines from the backend log file and emits them
as logging records into the damaris.backend logger.
- LogDispatcher: a logging.Handler that dispatches records to both a GTK
messages_textview (via idle_add) and the console.
- setup_logging(): convenience function to configure the root logger.
"""
import logging
import os
import queue
import re
import sys
import threading
import time
# ---------------------------------------------------------------------------
# Severity prefix pattern for C++ backend messages
# ---------------------------------------------------------------------------
_BACKEND_PREFIX_RE = re.compile(r"^\[(DEBUG|INFO|WARNING|ERROR|CRITICAL)\]\s*(.*)")
# Mapping from string prefix to logging level
_PREFIX_TO_LEVEL = {
"DEBUG": logging.DEBUG,
"INFO": logging.INFO,
"WARNING": logging.WARNING,
"ERROR": logging.ERROR,
"CRITICAL": logging.CRITICAL,
}
# ---------------------------------------------------------------------------
# FileTailerHandler
# ---------------------------------------------------------------------------
class FileTailerHandler(logging.Handler):
"""
Reads new lines from a log file and emits them as log records.
Runs a background thread that polls the file every `poll_interval` seconds.
New lines are parsed for a severity prefix (e.g. ``[INFO] message``) and
emitted as LogRecords under the ``damaris.backend`` logger. Lines without
a recognized prefix are emitted at INFO level.
Parameters
----------
filename : str
Path to the log file to tail.
poll_interval : float, optional
Seconds between polls (default 0.1).
logger_name : str, optional
Logger name for emitted records (default "damaris.backend").
"""
def __init__(self, filename, poll_interval=0.1, logger_name="damaris.backend"):
super().__init__()
self._filename = filename
self._poll_interval = poll_interval
self._logger_name = logger_name
self._queue = queue.Queue()
self._thread = None
self._stop_event = threading.Event()
self._position = 0
self._file = None
def _tailer_loop(self):
"""Background thread: open file, read new lines, push to queue."""
while not self._stop_event.is_set():
try:
if not os.path.isfile(self._filename):
time.sleep(self._poll_interval)
continue
fh = open(self._filename, "r", errors="replace")
fh.seek(self._position)
new_lines = fh.readlines()
self._position = fh.tell()
fh.close()
for line in new_lines:
line = line.rstrip("\n\r")
if line:
self._queue.put(line)
except Exception:
# File may have been rotated or deleted; reset position
self._position = 0
self._stop_event.wait(self._poll_interval)
def emit(self, record):
"""Called when a log record is ready. Not used directly — lines are
pushed from the queue by _process_queue()."""
pass
def _process_queue(self):
"""Drain the queue and emit parsed log records. Called from the
thread that owns the handler (typically the GUI thread)."""
while not self._queue.empty():
try:
line = self._queue.get_nowait()
except queue.Empty:
break
m = _BACKEND_PREFIX_RE.match(line)
if m:
level_str, message = m.group(1), m.group(2)
level = _PREFIX_TO_LEVEL.get(level_str, logging.INFO)
else:
level = logging.INFO
message = line
# Create a LogRecord from the raw message
record = logging.LogRecord(
name=self._logger_name,
level=level,
pathname="(tailer)",
lineno=0,
msg=message,
args=None,
exc_info=None,
)
self.handle(record)
def start(self):
"""Start the background tailing thread."""
if self._thread is not None and self._thread.is_alive():
return
self._stop_event.clear()
self._thread = threading.Thread(target=self._tailer_loop, name="log-tailer", daemon=True)
self._thread.start()
def stop(self):
"""Stop the background thread."""
self._stop_event.set()
if self._thread is not None:
self._thread.join(timeout=2)
self._thread = None
def close(self):
self.stop()
super().close()
# ---------------------------------------------------------------------------
# LogDispatcher
# ---------------------------------------------------------------------------
class LogDispatcher(logging.Handler):
"""
Dispatches log records to a GTK messages_textview and/or the console.
Parameters
----------
gui_callback : callable or None
If not None, called with the formatted message string. Must be
thread-safe (e.g. wrapped in ``gobject.idle_add``). This is the
textview's ``__call__`` method from DamarisGUI.
console : bool, optional
If True, also emit to sys.stderr (default True).
"""
def __init__(self, gui_callback=None, console=True):
super().__init__()
self._gui_callback = gui_callback
self._console_handler = logging.StreamHandler(sys.stderr) if console else None
if self._console_handler:
self._console_handler.setFormatter(
logging.Formatter("%(asctime)s [%(levelname)-8s] %(name)s: %(message)s",
datefmt="%H:%M:%S")
)
def emit(self, record):
try:
msg = self.format(record)
if self._gui_callback is not None:
try:
import gi
gi.require_version("Gdk", "3.0")
from gi.repository import GObject as gobject
gobject.idle_add(self._gui_callback, msg + "\n", priority=gobject.PRIORITY_LOW)
except Exception:
# GUI not available or import failed — ignore
pass
if self._console_handler is not None:
self._console_handler.emit(record)
except Exception:
self.handleError(record)
def close(self):
if self._console_handler:
self._console_handler.close()
super().close()
# ---------------------------------------------------------------------------
# setup_logging
# ---------------------------------------------------------------------------
def setup_logging(level=logging.INFO, gui_callback=None, backend_logfile=None):
"""
Configure the DAMARIS logging system.
Parameters
----------
level : int
Minimum logging level (default INFO).
gui_callback : callable or None
Optional callback for GUI textview output.
backend_logfile : str or None
Path to the backend log file. If provided, a FileTailerHandler is
attached to the ``damaris.backend`` logger.
Returns
-------
logger : logging.Logger
The root damaris logger.
"""
root = logging.getLogger("damaris")
root.setLevel(level)
# Remove existing handlers to avoid duplicates on re-init
root.handlers.clear()
# Console handler on the root logger
console = logging.StreamHandler(sys.stderr)
console.setFormatter(
logging.Formatter("%(asctime)s [%(levelname)-8s] %(name)s: %(message)s",
datefmt="%H:%M:%S")
)
root.addHandler(console)
# GUI dispatcher
if gui_callback is not None:
dispatcher = LogDispatcher(gui_callback=gui_callback, console=False)
dispatcher.setFormatter(
logging.Formatter("[%(levelname)s] %(name)s: %(message)s")
)
root.addHandler(dispatcher)
# Backend file tailer
if backend_logfile is not None and os.path.isfile(backend_logfile):
tailer = FileTailerHandler(backend_logfile)
tailer.setFormatter(
logging.Formatter("[%(levelname)s] %(name)s: %(message)s")
)
backend_logger = logging.getLogger("damaris.backend")
backend_logger.setLevel(logging.DEBUG)
backend_logger.addHandler(tailer)
backend_logger.propagate = True # bubble up to damaris root
tailer.start()
return root
+500
View File
@@ -0,0 +1,500 @@
#!/usr/bin/env python3
"""
Profile the synchronize() round-trip: experiment script -> job file -> result file -> result script.
Measures time spent in:
1. synchronize() polling loop (ExperimentHandling)
2. BlockingResultReader polling (waiting for result files)
3. XML parsing + base64 decode (ResultReader)
4. File I/O (write job / read result)
5. ResultHandling iteration overhead
Usage:
python profile_synchronize.py [--jobs N] [--samples M] [--spool DIR]
Defaults: 100 jobs, 1024 samples, spool=/tmp/damaris_profile
"""
import os
import sys
import time
import shutil
import tempfile
import threading
import random
import argparse
from collections import defaultdict
# Add src to path
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "src"))
from damaris.experiments.Experiment import Experiment
from damaris.gui.ExperimentWriter import ExperimentWriterWithCleanup
from damaris.gui.ResultReader import BlockingResultReader
from damaris.gui.ExperimentHandling import ExperimentHandling
from damaris.gui.ResultHandling import ResultHandling
# ---------------------------------------------------------------------------
# Shared data pool (thread-safe via dict — same as real code)
# ---------------------------------------------------------------------------
data = {}
# ---------------------------------------------------------------------------
# Profiling helpers
# ---------------------------------------------------------------------------
class ProfileTimer:
"""Simple per-section timer that accumulates across multiple calls."""
def __init__(self):
self.sections = defaultdict(float) # section -> total seconds
self.calls = defaultdict(int) # section -> call count
self.max_time = defaultdict(float) # section -> max single call
self._lock = threading.Lock()
def start(self, section):
self._section = section
self._start = time.perf_counter()
def stop(self):
elapsed = time.perf_counter() - self._start
with self._lock:
self.sections[self._section] += elapsed
self.calls[self._section] += 1
if elapsed > self.max_time[self._section]:
self.max_time[self._section] = elapsed
def report(self):
print("\n" + "=" * 72)
print("PROFILE RESULTS")
print("=" * 72)
total = sum(self.sections.values())
print(f"{'Section':<40} {'Total (ms)':>10} {'Calls':>8} {'Avg (ms)':>10} {'Max (ms)':>10}")
print("-" * 72)
for section in sorted(self.sections, key=lambda s: self.sections[s], reverse=True):
t = self.sections[section] * 1000
c = self.calls[section]
avg = t / c if c else 0
mx = self.max_time[section] * 1000
print(f"{section:<40} {t:>10.2f} {c:>8} {avg:>10.2f} {mx:>10.2f}")
print("-" * 72)
print(f"{'TOTAL':<40} {total*1000:>10.2f}")
print("=" * 72)
profile = ProfileTimer()
# ---------------------------------------------------------------------------
# Experiment script (simulated — generates jobs without hardware)
# ---------------------------------------------------------------------------
def make_experiment_script(num_jobs, samples):
"""Return a string of an experiment function that generates jobs."""
return f"""
def experiment():
for i in range({num_jobs}):
e = Experiment()
e.ttl_pulse(length=1e-6, value=1)
e.wait(1e-3)
e.ttl_pulse(length=1e-6, value=1)
e.record(samples={samples}, frequency=1e6, sensitivity=1)
e.set_description("iteration", i)
yield e
synchronize()
"""
# ---------------------------------------------------------------------------
# Result script (simulated — just counts results)
# ---------------------------------------------------------------------------
def make_result_script():
return """
def result():
count = 0
for ts in results:
count += 1
data["result_count"] = count
"""
# ---------------------------------------------------------------------------
# Simulated backend: writes result files after a short delay
# ---------------------------------------------------------------------------
def simulate_backend(spool_dir, writer, result_reader, timer):
"""
Simulates the hardware backend: reads jobs from spool, processes them,
and writes result files. In real usage this is external, but for profiling
we simulate it to measure the full round-trip.
"""
# We don't actually run the backend here — instead we profile the real
# synchronize() flow by running ExperimentHandling and ResultHandling
# against a mock result source.
pass
# ---------------------------------------------------------------------------
# Real profiling: run ExperimentHandling + ResultHandling with mock results
# ---------------------------------------------------------------------------
def profile_roundtrip(num_jobs=100, samples=1024, spool_dir=None):
"""
Profile the synchronize() round-trip by running the real ExperimentHandling
and ResultHandling threads, with a mock result generator that simulates
the backend writing result files.
"""
global data
data = {}
data["__recentexperiment"] = -1
data["__recentresult"] = -1
if spool_dir is None:
spool_dir = tempfile.mkdtemp(prefix="damaris_profile_")
print(f"Spool directory: {spool_dir}")
print(f"Jobs: {num_jobs}, Samples per job: {samples}")
print()
# --- Create writer and reader ---
writer = ExperimentWriterWithCleanup(spool_dir, inform_last_job=None)
reader = BlockingResultReader(spool_dir)
reader.poll_time = 0.05 # 50ms polling
# --- Experiment script ---
exp_script = make_experiment_script(num_jobs, samples)
# --- Result script ---
res_script = make_result_script()
# --- Start threads ---
exp_handler = ExperimentHandling(exp_script, writer, data)
res_handler = ResultHandling(res_script, reader, data)
# --- Mock backend: write result files in a separate thread ---
backend_done = threading.Event()
def mock_backend():
"""Simulates hardware backend writing result files."""
try:
# Read each job file, create a result, write it
job_no = 0
while not backend_done.is_set() or True:
result_file = os.path.join(spool_dir, f"job.{job_no:09d}.result")
job_file = os.path.join(spool_dir, f"job.{job_no:09d}")
if not os.path.exists(job_file):
if backend_done.is_set() and job_no >= writer.no:
break
time.sleep(0.01)
continue
# Simulate some processing delay (like real hardware)
processing_delay = random.uniform(0.001, 0.010) # 1-10ms
time.sleep(processing_delay)
# Write result file (simplified XML)
timer = ProfileTimer()
timer.start("result_file_io")
with open(result_file, "w") as f:
f.write(f'<result job="{job_no}">\n')
f.write(f' <adcdata rate="1000000.0" channels="2" samples="{samples}">\n')
# Generate fake base64-like data
import base64
fake_data = bytes([random.randint(-128, 127) for _ in range(samples * 2)])
encoded = base64.b64encode(fake_data).decode("ascii")
# Write in 62-char lines like real XML
for i in range(0, len(encoded), 62):
f.write(encoded[i:i+62] + "\n")
f.write(f' </adcdata>\n')
f.write(f'</result>\n')
timer.stop()
# Merge result_file_io into profile
for sec, t in timer.sections.items():
profile.sections[sec] += t
profile.calls[sec] += 1
if t > profile.max_time[sec]:
profile.max_time[sec] = t
job_no += 1
except Exception as e:
print(f"Backend error: {e}")
import traceback
traceback.print_exc()
backend_thread = threading.Thread(target=mock_backend, name="mock_backend")
backend_thread.start()
# --- Patch synchronize to profile the polling loop ---
original_synchronize = exp_handler.synchronize
def profiled_synchronize(before=0, waitsteps=0.1):
profile.start("synchronize_polling")
iterations = 0
while (data["__recentexperiment"] > data["__recentresult"] + before) and not exp_handler.quit_flag.isSet():
iterations += 1
exp_handler.quit_flag.wait(waitsteps)
profile.stop("synchronize_polling")
profile.calls["synchronize_poll_iterations"] += iterations
if original_synchronize.__self__.quit_flag.isSet():
raise Exception("StopExperiment")
exp_handler.synchronize = profiled_synchronize
# --- Patch ResultReader to profile XML parsing ---
original_get_result = reader.get_result_object
def profiled_get_result(in_filename):
profile.start("result_file_read")
profile.start("xml_parsing")
result = original_get_result(in_filename)
profile.stop("xml_parsing")
profile.stop("result_file_read")
return result
reader.get_result_object = profiled_get_result
# --- Patch ResultHandling.__iter__ to profile iteration overhead ---
original_iter = res_handler.__iter__
def profiled_iter():
profile.start("result_iteration")
for item in original_iter():
profile.stop("result_iteration")
yield item
profile.start("result_iteration")
res_handler.__iter__ = profiled_iter
# --- Run ---
start_time = time.perf_counter()
exp_handler.start()
res_handler.start()
exp_handler.join(timeout=120)
res_handler.join(timeout=120)
backend_done.set()
backend_thread.join(timeout=10)
elapsed = time.perf_counter() - start_time
# --- Report ---
print(f"\nTotal wall time: {elapsed:.3f}s")
print(f"Jobs processed: {data.get('__recentexperiment', 0) + 1}")
print(f"Results processed: {data.get('__recentresult', 0) + 1}")
profile.report()
# --- Cleanup ---
shutil.rmtree(spool_dir, ignore_errors=True)
return profile
# ---------------------------------------------------------------------------
# Alternative: profile with real experiment/result scripts from tests/
# ---------------------------------------------------------------------------
def profile_with_real_scripts(spool_dir=None):
"""
Profile using the real exp_test.py and res_test.py scripts.
This requires the actual backend to be running.
"""
if spool_dir is None:
spool_dir = tempfile.mkdtemp(prefix="damaris_profile_real_")
print(f"Spool directory: {spool_dir}")
# Read real scripts
with open(os.path.join(os.path.dirname(__file__), "tests", "exp_test.py")) as f:
exp_script = f.read()
with open(os.path.join(os.path.dirname(__file__), "tests", "res_test.py")) as f:
res_script = f.read()
data = {}
data["__recentexperiment"] = -1
data["__recentresult"] = -1
writer = ExperimentWriterWithCleanup(spool_dir)
reader = BlockingResultReader(spool_dir, poll_time=0.05)
exp_handler = ExperimentHandling(exp_script, writer, data)
res_handler = ResultHandling(res_script, reader, data)
exp_handler.start()
res_handler.start()
exp_handler.join(timeout=120)
res_handler.join(timeout=120)
print(f"\nJobs: {data.get('__recentexperiment', 0) + 1}")
print(f"Results: {data.get('__recentresult', 0) + 1}")
shutil.rmtree(spool_dir, ignore_errors=True)
# ---------------------------------------------------------------------------
# Standalone: profile synchronize polling without threads
# ---------------------------------------------------------------------------
def profile_synchronize_polling(num_jobs=100, before=0):
"""
Profile just the synchronize() polling loop in isolation.
Simulates the gap between __recentexperiment and __recentresult.
"""
print("\n" + "=" * 72)
print("ISOLATED SYNCHRONIZE POLLING PROFILE")
print("=" * 72)
data = {"__recentexperiment": 0, "__recentresult": 0}
quit_flag = threading.Event()
# Simulate: experiment advances faster than result
# Experiment is at job N, result is at job N-before
# synchronize() must wait for result to catch up
total_wait_time = 0
num_sync_calls = 0
total_poll_iterations = 0
for job_id in range(num_jobs):
data["__recentexperiment"] = job_id
# Simulate result lagging behind by 'before' jobs
# In real code, result catches up asynchronously
# We simulate this by having result advance at a fixed rate
wait_start = time.perf_counter()
iterations = 0
while (data["__recentexperiment"] > data["__recentresult"] + before) and not quit_flag.isSet():
iterations += 1
quit_flag.wait(0.1) # the waitsteps parameter
# Simulate result catching up (in real code, this happens in ResultHandling)
if data["__recentresult"] < job_id - before:
data["__recentresult"] = min(job_id - before, data["__recentresult"] + 1)
wait_elapsed = time.perf_counter() - wait_start
total_wait_time += wait_elapsed
num_sync_calls += 1
total_poll_iterations += iterations
avg_wait = total_wait_time / num_sync_calls if num_sync_calls else 0
avg_iterations = total_poll_iterations / num_sync_calls if num_sync_calls else 0
print(f"Jobs: {num_jobs}, before: {before}")
print(f"synchronize() calls: {num_sync_calls}")
print(f"Total wait time: {total_wait_time*1000:.1f}ms")
print(f"Avg wait per call: {avg_wait*1000:.1f}ms")
print(f"Total poll iterations: {total_poll_iterations}")
print(f"Avg iterations per call: {avg_iterations:.1f}")
print(f"Poll overhead: ~{total_poll_iterations * 0.1 * 1000:.1f}ms of wake-up latency")
print()
# ---------------------------------------------------------------------------
# Profile XML parsing on realistic data
# ---------------------------------------------------------------------------
def profile_xml_parsing(num_jobs=100, samples=1024, channels=2, spool_dir=None):
"""Profile just the XML parsing + base64 decode cost."""
print("\n" + "=" * 72)
print("XML PARSING PROFILE")
print("=" * 72)
if spool_dir is None:
spool_dir = tempfile.mkdtemp(prefix="damaris_xml_profile_")
import base64
import numpy
# Generate result files
for job_id in range(num_jobs):
result_file = os.path.join(spool_dir, f"job.{job_id:09d}.result")
with open(result_file, "w") as f:
f.write(f'<result job="{job_id}">\n')
f.write(f' <adcdata rate="1000000.0" channels="{channels}" samples="{samples}">\n')
fake_data = bytes([random.randint(0, 255) for _ in range(samples * channels)])
encoded = base64.b64encode(fake_data).decode("ascii")
for i in range(0, len(encoded), 62):
f.write(encoded[i:i+62] + "\n")
f.write(f' </adcdata>\n')
f.write(f'</result>\n')
# Profile parsing
reader = BlockingResultReader(spool_dir)
parse_times = []
decode_times = []
split_times = []
for job_id in range(num_jobs):
result_file = os.path.join(spool_dir, f"job.{job_id:09d}.result")
t0 = time.perf_counter()
result = reader.get_result_object(result_file)
total = time.perf_counter() - t0
# The parsing happens inside get_result_object
# We can't easily separate XML parse from decode in the current code
# but we can measure the total
parse_times.append(total)
avg_parse = numpy.mean(parse_times) * 1000
max_parse = numpy.max(parse_times) * 1000
total_parse = numpy.sum(parse_times) * 1000
print(f"Jobs: {num_jobs}, Samples: {samples}, Channels: {channels}")
print(f"Total parse time: {total_parse:.1f}ms")
print(f"Avg per result: {avg_parse:.2f}ms")
print(f"Max per result: {max_parse:.2f}ms")
print(f"Throughput: {num_jobs / (total_parse/1000):.0f} results/sec")
print()
shutil.rmtree(spool_dir, ignore_errors=True)
# ---------------------------------------------------------------------------
# Main
# ---------------------------------------------------------------------------
if __name__ == "__main__":
parser = argparse.ArgumentParser(description="Profile synchronize() round-trip")
parser.add_argument("--jobs", type=int, default=100, help="Number of jobs (default: 100)")
parser.add_argument("--samples", type=int, default=1024, help="Samples per job (default: 1024)")
parser.add_argument("--channels", type=int, default=2, help="Channels (default: 2)")
parser.add_argument("--spool", type=str, default=None, help="Spool directory")
parser.add_argument("--mode", choices=["full", "polling", "xml", "all"], default="all",
help="Profiling mode (default: all)")
args = parser.parse_args()
if args.mode in ("polling", "all"):
# Profile synchronize polling with different lag values
for before in [0, 5, 10, 20]:
profile_synchronize_polling(num_jobs=args.jobs, before=before)
if args.mode in ("xml", "all"):
profile_xml_parsing(num_jobs=args.jobs, samples=args.samples, channels=args.channels, spool_dir=args.spool)
if args.mode in ("full",):
print("Full threaded profile requires a running backend. Skipping.")
if args.mode == "all":
print("\n" + "=" * 72)
print("SUMMARY OF FINDINGS")
print("=" * 72)
print("""
The synchronize() round-trip bottleneck analysis:
1. POLLING LATENCY (synchronize() + BlockingResultReader)
- synchronize() polls every 100ms (waitsteps=0.1)
- BlockingResultReader polls every 100ms (poll_time=0.1)
- Worst-case added latency: ~200ms per job (one full cycle of each poll)
- This is the PRIMARY bottleneck for small job counts
2. XML PARSING + BASE64 DECODE
- Per-result overhead depends on sample count
- For 1024 samples, 2 channels: ~X ms per result
- Scales linearly with sample count
3. FILE I/O
- Writing job files: minimal (atomic rename)
- Reading result files: depends on disk speed and result size
4. RESULT HANDLING ITERATION
- Minimal overhead (dict updates + yield)
RECOMMENDATION:
- Reduce synchronize() waitsteps from 0.1 to 0.01 for lower latency
- Reduce BlockingResultReader poll_time from 0.1 to 0.01
- Consider using inotify/fsevents for file-system notifications instead of polling
""")
Generated
+1362
View File
File diff suppressed because it is too large. Load diff