Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
0c476c6
Initial commit
Koldstream Feb 26, 2026
794ac09
Moving in latest main
Koldstream Feb 26, 2026
ffa5651
Merge in latest main
Koldstream Mar 5, 2026
8870254
Commenting this out for now
Koldstream Mar 5, 2026
f19137e
Removing static imports
Koldstream Apr 7, 2026
50b7aa7
Updating job number handlers
Koldstream Apr 15, 2026
5be9a35
Suppress unwanted warnings
stephen-riggs Apr 16, 2026
e8c41e4
Merge branch 'main' into doppio-live-processing
stephen-riggs Apr 16, 2026
5a46af1
Fix logging
stephen-riggs Apr 16, 2026
356cd94
Restore default values
stephen-riggs Apr 27, 2026
35ca6f7
Remove picker id
stephen-riggs Apr 28, 2026
82854eb
Merging main
Koldstream Jun 10, 2026
49028be
Fixing job number reservations, sqlite db locking, minor db bugs, and…
Koldstream Jul 3, 2026
475221e
Merge branch 'main' into doppio-live-processing
stephen-riggs Jul 7, 2026
f43be58
Bits got removed inadvertantly
stephen-riggs Jul 9, 2026
8e1694a
Pipeliner as dependency, and try to make codeql happy
stephen-riggs Jul 13, 2026
d61c10c
Merge branch 'main' into doppio-live-processing
stephen-riggs Jul 17, 2026
cceec72
Update mocked transport object in tests
stephen-riggs Jul 17, 2026
b981dca
Standardise transport object import
stephen-riggs Jul 17, 2026
98ce0fa
Update more tests for transport object
stephen-riggs Jul 17, 2026
4adc612
Merge branch 'main' into doppio-live-processing
stephen-riggs Aug 14, 2026
afcb295
post-merge cleanup
stephen-riggs Aug 14, 2026
2b843a7
Restore default murfey behaviour for mc mrc
stephen-riggs Aug 14, 2026
e2554fb
2D and refine jobs for symmetry and edge cases
stephen-riggs Aug 14, 2026
00d1dfd
Attempt to restore non-doppio job number behaviour
stephen-riggs Aug 14, 2026
de051cc
Update tests for new behaviour of reserving
stephen-riggs Aug 14, 2026
b821126
Remove the other_options which duplicates feedback params
stephen-riggs Aug 14, 2026
b737828
Correct type hinting
stephen-riggs Aug 17, 2026
ce58001
Correct type hinting
stephen-riggs Aug 17, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@ developer = [
"pytest-mock", # Additional mocking tools for unit tests
]
server = [
"ccpem-pipeliner",
"cryptography",
"graypy",
"ispyb>=12.1.0", # Responsible for setting requirements for SQLAlchemy and mysql-connector-python;
Expand Down
14 changes: 4 additions & 10 deletions src/murfey/cli/inject_spa_processing.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,6 @@
from murfey.util.config import get_machine_config, get_microscope, get_security_config
from murfey.util.db import (
AutoProcProgram,
ClassificationFeedbackParameters,
ClientEnvironment,
DataCollection,
DataCollectionGroup,
Expand Down Expand Up @@ -136,15 +135,11 @@ def run():
.where(AutoProcProgram.pj_id == ProcessingJob.id)
.where(ProcessingJob.recipe == "em-spa-preprocess")
).one()
params = murfey_db.exec(
select(SPARelionParameters, ClassificationFeedbackParameters)
.where(SPARelionParameters.pj_id == collected_ids[2].id)
.where(ClassificationFeedbackParameters.pj_id == SPARelionParameters.pj_id)
proc_params = murfey_db.exec(
select(SPARelionParameters).where(
SPARelionParameters.pj_id == collected_ids[2].id
)
).one()
proc_params: dict | None = dict(params[0])
feedback_params = params[1]
if feedback_params.picker_murfey_id is None:
raise ValueError("No ISPyB picker ID was found")
except sqlalchemy.exc.NoResultFound:
proc_params = None

Expand Down Expand Up @@ -196,7 +191,6 @@ def run():
"ft_bin": proc_params["motion_corr_binning"],
"fm_dose": proc_params["dose_per_frame"],
"gain_ref": proc_params["gain_ref"],
"picker_uuid": feedback_params.picker_murfey_id,
"session_id": args.session_id,
"particle_diameter": proc_params["particle_diameter"] or 0,
"fm_int_file": args.eer_fractionation_file,
Expand Down
5 changes: 5 additions & 0 deletions src/murfey/server/api/auth.py
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,8 @@ async def submit_to_auth_endpoint(
Helper function to forward incoming requests to an authentication server
to verify that they are allowed to inspect the
"""
if security_config.auth_type == "none":
return {"valid": True}

# Forward only essentials auth-related headers
headers = {
Expand Down Expand Up @@ -189,6 +191,9 @@ async def validate_instrument_token(
"""
Used by the backend routers to check the incoming instrument server token.
"""
if security_config.instrument_auth_type == "none":
return None

try:
# Validate using auth URL if provided
if security_config.instrument_auth_url:
Expand Down
12 changes: 7 additions & 5 deletions src/murfey/server/api/session_control.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,8 +25,8 @@
keycloak_client = None
SMARTEM_ACTIVE = False

import murfey.server
import murfey.server.prometheus as prom
from murfey.server import _transport_object
from murfey.server.api.auth import (
MurfeySessionIDInstrument as MurfeySessionID,
validate_instrument_token,
Expand Down Expand Up @@ -200,14 +200,14 @@ def register_processing_success_in_ispyb(
.where(AutoProcProgram.pj_id == ProcessingJob.id)
).all()
appids = [c[3].id for c in collected_ids]
if _transport_object:
if murfey.server._transport_object:
if db is not None:
apps = db.query(ISPyBAutoProcProgram).filter(
ISPyBAutoProcProgram.autoProcProgramId.in_(appids)
)
for updated in apps:
updated.processingStatus = True
_transport_object.do_update_processing_status(updated)
murfey.server._transport_object.do_update_processing_status(updated)


@router.get("/num_movies")
Expand All @@ -234,8 +234,10 @@ def failed_client_post(instrument_name: str, post_info: PostInfo):
"data": post_info.data,
"kwargs": post_info.kwargs,
}
if _transport_object:
_transport_object.send(_transport_object.feedback_queue, zocalo_message)
if murfey.server._transport_object:
murfey.server._transport_object.send(
murfey.server._transport_object.feedback_queue, zocalo_message
)


@router.post("/sessions/{session_id}/rsyncer")
Expand Down
94 changes: 54 additions & 40 deletions src/murfey/server/api/session_info.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,11 +7,11 @@
from fastapi import APIRouter, Depends, Request
from fastapi.responses import FileResponse, JSONResponse
from pydantic import BaseModel
from sqlmodel import select
from sqlmodel import Session as SQLModelSession, select

import murfey.server
import murfey.server.api.websocket as ws
import murfey.server.prometheus as prom
from murfey.server import _transport_object
from murfey.server.api import templates
from murfey.server.api.auth import (
MurfeyInstrumentNameFrontend as MurfeyInstrumentName,
Expand Down Expand Up @@ -44,7 +44,7 @@
Movie,
ProcessingJob,
RsyncInstance,
Session,
Session as MurfeySession,
SessionProcessingParameters,
SPARelionParameters,
Tilt,
Expand All @@ -67,7 +67,7 @@ def health_check(db=ispyb_db):
conn.close()
return {
"ispyb_connection": True,
"rabbitmq_connection": _transport_object.transport.is_connected(),
"rabbitmq_connection": murfey.server._transport_object.transport.is_connected(),
}


Expand Down Expand Up @@ -131,30 +131,34 @@ def all_visit_info(


@router.get("/sessions/{session_id}/rsyncers", response_model=List[RsyncInstance])
def get_rsyncers_for_client(session_id: MurfeySessionID, db=murfey_db):
def get_rsyncers_for_client(
session_id: MurfeySessionID, db: SQLModelSession = murfey_db
):
rsync_instances = db.exec(
select(RsyncInstance).where(RsyncInstance.session_id == session_id)
)
return rsync_instances.all()


class SessionClients(BaseModel):
session: Session
session: MurfeySession
clients: List[ClientEnvironment]


@router.get("/sessions/{session_id}")
async def get_session(session_id: MurfeySessionID, db=murfey_db) -> SessionClients:
session = db.exec(select(Session).where(Session.id == session_id)).one()
async def get_session(
session_id: MurfeySessionID, db: SQLModelSession = murfey_db
) -> SessionClients:
session = db.exec(select(MurfeySession).where(MurfeySession.id == session_id)).one()
clients = db.exec(
select(ClientEnvironment).where(ClientEnvironment.session_id == session_id)
).all()
return SessionClients(session=session, clients=clients)


@router.get("/sessions")
async def get_sessions(db=murfey_db):
sessions = db.exec(select(Session)).all()
async def get_sessions(db: SQLModelSession = murfey_db):
sessions = db.exec(select(MurfeySession)).all()
clients = db.exec(select(ClientEnvironment)).all()
res = []
for sess in sessions:
Expand All @@ -176,9 +180,9 @@ def create_session(
visit: str,
name: str,
visit_end_time: VisitEndTime,
db=murfey_db,
db: SQLModelSession = murfey_db,
) -> int:
s = Session(
s = MurfeySession(
name=name,
visit=visit,
instrument_name=instrument_name,
Expand All @@ -199,9 +203,9 @@ def update_session(
session_id: MurfeySessionID,
process: bool = True,
smartem_acquisition_uuid: str | None = None,
db=murfey_db,
db: SQLModelSession = murfey_db,
) -> None:
session = db.exec(select(Session).where(Session.id == session_id)).one()
session = db.exec(select(MurfeySession).where(MurfeySession.id == session_id)).one()
session.process = process
session.smartem_acquisition_uuid = smartem_acquisition_uuid
db.add(session)
Expand All @@ -210,35 +214,37 @@ def update_session(


@router.delete("/sessions/{session_id}")
def remove_session(session_id: MurfeySessionID, db=murfey_db):
def remove_session(session_id: MurfeySessionID, db: SQLModelSession = murfey_db):
remove_session_by_id(session_id, db)


@router.get("/instruments/{instrument_name}/visits/{visit_name}/sessions")
def get_sessions_with_visit(
instrument_name: MurfeyInstrumentName, visit_name: str, db=murfey_db
) -> List[Session]:
instrument_name: MurfeyInstrumentName,
visit_name: str,
db: SQLModelSession = murfey_db,
) -> List[MurfeySession]:
sessions = db.exec(
select(Session)
.where(Session.instrument_name == instrument_name)
.where(Session.visit == visit_name)
select(MurfeySession)
.where(MurfeySession.instrument_name == instrument_name)
.where(MurfeySession.visit == visit_name)
).all()
return sessions


@router.get("/instruments/{instrument_name}/sessions")
async def get_sessions_by_instrument_name(
instrument_name: MurfeyInstrumentName, db=murfey_db
) -> List[Session]:
instrument_name: MurfeyInstrumentName, db: SQLModelSession = murfey_db
) -> List[MurfeySession]:
sessions = db.exec(
select(Session).where(Session.instrument_name == instrument_name)
select(MurfeySession).where(MurfeySession.instrument_name == instrument_name)
).all()
return sessions


@router.get("/sessions/{session_id}/data_collection_groups")
def get_dc_groups(
session_id: MurfeySessionID, db=murfey_db
session_id: MurfeySessionID, db: SQLModelSession = murfey_db
) -> Dict[str, DataCollectionGroup]:
data_collection_groups = db.exec(
select(DataCollectionGroup).where(DataCollectionGroup.session_id == session_id)
Expand All @@ -248,7 +254,7 @@ def get_dc_groups(

@router.get("/sessions/{session_id}/data_collection_groups/{dcgid}/data_collections")
def get_data_collections(
session_id: MurfeySessionID, dcgid: int, db=murfey_db
session_id: MurfeySessionID, dcgid: int, db: SQLModelSession = murfey_db
) -> List[DataCollection]:
data_collections = db.exec(
select(DataCollection).where(DataCollection.dcg_id == dcgid)
Expand All @@ -257,7 +263,7 @@ def get_data_collections(


@router.get("/clients")
async def get_clients(db=murfey_db):
async def get_clients(db: SQLModelSession = murfey_db):
clients = db.exec(select(ClientEnvironment)).all()
return clients

Expand All @@ -268,9 +274,11 @@ class CurrentGainRef(BaseModel):

@router.put("/sessions/{session_id}/current_gain_ref")
def update_current_gain_ref(
session_id: MurfeySessionID, new_gain_ref: CurrentGainRef, db=murfey_db
session_id: MurfeySessionID,
new_gain_ref: CurrentGainRef,
db: SQLModelSession = murfey_db,
):
session = db.exec(select(Session).where(Session.id == session_id)).one()
session = db.exec(select(MurfeySession).where(MurfeySession.id == session_id)).one()
session.current_gain_ref = new_gain_ref.path
db.add(session)

Expand Down Expand Up @@ -391,7 +399,7 @@ class ProcessingDetails(BaseModel):

@spa_router.get("/sessions/{session_id}/spa_processing_parameters")
def get_spa_proc_param_details(
session_id: MurfeySessionID, db=murfey_db
session_id: MurfeySessionID, db: SQLModelSession = murfey_db
) -> Optional[List[ProcessingDetails]]:
params = db.exec(
select(
Expand Down Expand Up @@ -440,7 +448,7 @@ def _parse(ps, i, dcg_id):
"/sessions/{session_id}/data_collection_groups/{dcgid}/grid_squares/{gsid}/foil_holes/{fhid}/num_movies"
)
def get_number_of_movies_from_foil_hole(
session_id: int, dcgid: int, gsid: int, fhid: int, db=murfey_db
session_id: int, dcgid: int, gsid: int, fhid: int, db: SQLModelSession = murfey_db
) -> int:
movies = db.exec(
select(Movie, FoilHole, GridSquare, DataCollectionGroup)
Expand All @@ -456,13 +464,13 @@ def get_number_of_movies_from_foil_hole(


@spa_router.get("/sessions/{session_id}/grid_squares")
def get_grid_squares(session_id: MurfeySessionID, db=murfey_db):
def get_grid_squares(session_id: MurfeySessionID, db: SQLModelSession = murfey_db):
return _get_grid_squares(session_id, db)


@spa_router.get("/sessions/{session_id}/data_collection_groups/{dcgid}/grid_squares")
def get_grid_squares_from_dcg(
session_id: MurfeySessionID, dcgid: int, db=murfey_db
session_id: MurfeySessionID, dcgid: int, db: SQLModelSession = murfey_db
) -> List[GridSquare]:
return _get_grid_squares_from_dcg(session_id, dcgid, db)

Expand All @@ -471,14 +479,14 @@ def get_grid_squares_from_dcg(
"/sessions/{session_id}/data_collection_groups/{dcgid}/grid_squares/{gsid}/foil_holes"
)
def get_foil_holes_from_grid_square(
session_id: MurfeySessionID, dcgid: int, gsid: int, db=murfey_db
session_id: MurfeySessionID, dcgid: int, gsid: int, db: SQLModelSession = murfey_db
) -> List[FoilHole]:
return _get_foil_holes_from_grid_square(session_id, dcgid, gsid, db)


@spa_router.get("/sessions/{session_id}/foil_hole/{fh_name}")
def get_foil_hole(
session_id: MurfeySessionID, fh_name: int, db=murfey_db
session_id: MurfeySessionID, fh_name: int, db: SQLModelSession = murfey_db
) -> Dict[str, int]:
return _get_foil_hole(session_id, fh_name, db)

Expand All @@ -492,7 +500,7 @@ def get_foil_hole(

@tomo_router.get("/sessions/{session_id}/tilt_series/{tilt_series_tag}/tilts")
def get_tilts(
session_id: MurfeySessionID, tilt_series_tag: str, db=murfey_db
session_id: MurfeySessionID, tilt_series_tag: str, db: SQLModelSession = murfey_db
) -> Dict[str, List[str]]:
res = db.exec(
select(TiltSeries, Tilt)
Expand All @@ -517,7 +525,9 @@ def get_tilts(


@correlative_router.get("/sessions/{session_id}/upstream_visits")
async def find_upstream_visits(session_id: MurfeySessionID, db=murfey_db):
async def find_upstream_visits(
session_id: MurfeySessionID, db: SQLModelSession = murfey_db
):
return _find_upstream_visits(session_id=session_id, db=db)


Expand All @@ -528,7 +538,7 @@ async def gather_upstream_files(
visit_name: str,
session_id: MurfeySessionID,
upstream_file_request: UpstreamFileRequestInfo,
db=murfey_db,
db: SQLModelSession = murfey_db,
):
return _gather_upstream_files(
session_id=session_id,
Expand All @@ -546,7 +556,7 @@ async def get_upstream_file(
visit_name: str,
session_id: MurfeySessionID,
upstream_file_path: Path,
db=murfey_db,
db: SQLModelSession = murfey_db,
):
upstream_file = _get_upstream_file(upstream_file_path)
return (
Expand All @@ -557,14 +567,18 @@ async def get_upstream_file(
@correlative_router.get(
"/visits/{visit_name}/sessions/{session_id}/upstream_tiff_paths"
)
async def gather_upstream_tiffs(visit_name: str, session_id: int, db=murfey_db):
async def gather_upstream_tiffs(
visit_name: str, session_id: int, db: SQLModelSession = murfey_db
):
return _gather_upstream_tiffs(visit_name=visit_name, session_id=session_id, db=db)


@correlative_router.get(
"/visits/{visit_name}/sessions/{session_id}/upstream_tiff/{tiff_path:path}"
)
async def get_tiff_file(visit_name: str, session_id: int, tiff_path: str, db=murfey_db):
async def get_tiff_file(
visit_name: str, session_id: int, tiff_path: str, db: SQLModelSession = murfey_db
):
tiff_file = _get_tiff_file(
visit_name=visit_name, session_id=session_id, tiff_path=tiff_path, db=db
)
Expand Down
Loading