信息架构/菜单框架大重构-20260109

This commit is contained in:
Cheng Zhou
2026-01-09 11:19:04 +08:00
parent ba3cf95b8a
commit 9021c7fe2b
188 changed files with 7632 additions and 11658 deletions
+5
View File
@@ -113,3 +113,8 @@ async def update_ae(db: AsyncSession, study_id: uuid.UUID, ae: AdverseEvent, ae_
await db.commit()
await db.refresh(ae)
return ae
async def delete_ae(db: AsyncSession, ae: AdverseEvent) -> None:
await db.delete(ae)
await db.commit()
-81
View File
@@ -1,81 +0,0 @@
import uuid
from typing import Sequence
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.comment import Comment
from app.schemas.comment import CommentCreate
async def create_comment(
db: AsyncSession,
study_id: uuid.UUID,
entity_type: str,
entity_id: uuid.UUID,
comment_in: CommentCreate,
created_by: uuid.UUID,
) -> Comment:
comment = Comment(
study_id=study_id,
entity_type=entity_type,
entity_id=entity_id,
content=comment_in.content,
quote_comment_id=comment_in.quote_comment_id,
created_by=created_by,
)
db.add(comment)
await db.commit()
await db.refresh(comment)
return comment
async def list_comments(
db: AsyncSession,
study_id: uuid.UUID,
entity_type: str,
entity_id: uuid.UUID,
) -> Sequence[Comment]:
result = await db.execute(
select(Comment)
.where(
Comment.study_id == study_id,
Comment.entity_type == entity_type,
Comment.entity_id == entity_id,
Comment.is_deleted.is_(False),
)
.order_by(Comment.created_at.asc())
)
return result.scalars().all()
async def get_comment(db: AsyncSession, comment_id: uuid.UUID, *, include_deleted: bool = False) -> Comment | None:
stmt = select(Comment).where(Comment.id == comment_id)
if not include_deleted:
stmt = stmt.where(Comment.is_deleted.is_(False))
result = await db.execute(stmt)
return result.scalar_one_or_none()
async def get_comments_by_ids(
db: AsyncSession,
ids: set[uuid.UUID],
*,
include_deleted: bool = False,
) -> dict[uuid.UUID, Comment]:
if not ids:
return {}
stmt = select(Comment).where(Comment.id.in_(ids))
if not include_deleted:
stmt = stmt.where(Comment.is_deleted.is_(False))
result = await db.execute(stmt)
comments = result.scalars().all()
return {c.id: c for c in comments}
async def soft_delete_comment(db: AsyncSession, comment: Comment) -> Comment:
comment.is_deleted = True
db.add(comment)
await db.commit()
await db.refresh(comment)
return comment
-117
View File
@@ -1,117 +0,0 @@
import uuid
from datetime import date, datetime, timezone
from typing import Sequence
from sqlalchemy import select, update as sa_update
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.data_query import DataQuery
from app.models.site import Site
from app.models.subject import Subject
from app.schemas.data_query import DataQueryCreate, DataQueryUpdate
async def _validate_site_subject(db: AsyncSession, study_id: uuid.UUID, site_id: uuid.UUID | None, subject_id: uuid.UUID | None):
if site_id:
result = await db.execute(select(Site).where(Site.id == site_id))
site = result.scalar_one_or_none()
if not site or site.study_id != study_id:
raise ValueError("Site not found in study")
subject_site = None
if subject_id:
result = await db.execute(select(Subject).where(Subject.id == subject_id))
subj = result.scalar_one_or_none()
if not subj or subj.study_id != study_id:
raise ValueError("Subject not found in study")
subject_site = subj.site_id
if site_id and subject_site and site_id != subject_site:
raise ValueError("Site and subject mismatch")
async def create_query(
db: AsyncSession,
study_id: uuid.UUID,
query_in: DataQueryCreate,
*,
created_by: uuid.UUID,
) -> DataQuery:
await _validate_site_subject(db, study_id, query_in.site_id, query_in.subject_id)
dq = DataQuery(
study_id=study_id,
site_id=query_in.site_id,
subject_id=query_in.subject_id,
visit_id=query_in.visit_id,
title=query_in.title,
description=query_in.description,
category=query_in.category,
priority=query_in.priority,
assigned_to=query_in.assigned_to,
due_date=query_in.due_date,
status="OPEN",
resolution=None,
closed_at=None,
created_by=created_by,
)
db.add(dq)
await db.commit()
await db.refresh(dq)
return dq
async def get_query(db: AsyncSession, query_id: uuid.UUID) -> DataQuery | None:
result = await db.execute(select(DataQuery).where(DataQuery.id == query_id))
return result.scalar_one_or_none()
async def list_queries(
db: AsyncSession,
study_id: uuid.UUID,
*,
status: str | None = None,
site_id: uuid.UUID | None = None,
subject_id: uuid.UUID | None = None,
assigned_to: uuid.UUID | None = None,
overdue: bool | None = None,
category: str | None = None,
priority: str | None = None,
skip: int = 0,
limit: int = 100,
) -> Sequence[DataQuery]:
stmt = select(DataQuery).where(DataQuery.study_id == study_id)
if status:
stmt = stmt.where(DataQuery.status == status)
if site_id:
stmt = stmt.where(DataQuery.site_id == site_id)
if subject_id:
stmt = stmt.where(DataQuery.subject_id == subject_id)
if assigned_to:
stmt = stmt.where(DataQuery.assigned_to == assigned_to)
if category:
stmt = stmt.where(DataQuery.category == category)
if priority:
stmt = stmt.where(DataQuery.priority == priority)
if overdue is True:
stmt = stmt.where(DataQuery.due_date < date.today(), DataQuery.status != "CLOSED")
if overdue is False:
stmt = stmt.where((DataQuery.due_date >= date.today()) | (DataQuery.due_date.is_(None)) | (DataQuery.status == "CLOSED"))
stmt = stmt.offset(skip).limit(limit)
result = await db.execute(stmt)
return result.scalars().all()
async def update_query(db: AsyncSession, dq: DataQuery, dq_in: DataQueryUpdate) -> DataQuery:
update_data = dq_in.model_dump(exclude_unset=True)
if "status" in update_data:
if update_data["status"] == "CLOSED":
update_data["closed_at"] = datetime.now(timezone.utc)
else:
update_data["closed_at"] = None
if update_data:
await db.execute(
sa_update(DataQuery)
.where(DataQuery.id == dq.id)
.values(**update_data)
)
await db.commit()
await db.refresh(dq)
return dq
+76
View File
@@ -0,0 +1,76 @@
import uuid
from typing import Sequence
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.drug_shipment import DrugShipment
from app.schemas.drug_shipment import DrugShipmentCreate, DrugShipmentUpdate
async def create_shipment(
db: AsyncSession,
study_id: uuid.UUID,
shipment_in: DrugShipmentCreate,
created_by: uuid.UUID | None,
) -> DrugShipment:
shipment = DrugShipment(
study_id=study_id,
direction=shipment_in.direction,
site_name=shipment_in.site_name,
drug_desc=shipment_in.drug_desc,
ship_date=shipment_in.ship_date,
carrier=shipment_in.carrier,
tracking_no=shipment_in.tracking_no,
from_party=shipment_in.from_party,
to_party=shipment_in.to_party,
status=shipment_in.status,
remark=shipment_in.remark,
created_by=created_by,
)
db.add(shipment)
await db.commit()
await db.refresh(shipment)
return shipment
async def get_shipment(db: AsyncSession, shipment_id: uuid.UUID) -> DrugShipment | None:
result = await db.execute(select(DrugShipment).where(DrugShipment.id == shipment_id))
return result.scalar_one_or_none()
async def list_shipments(
db: AsyncSession,
study_id: uuid.UUID,
site_name: str | None = None,
tracking_no: str | None = None,
status: str | None = None,
skip: int = 0,
limit: int = 100,
) -> Sequence[DrugShipment]:
stmt = select(DrugShipment).where(DrugShipment.study_id == study_id)
if site_name:
stmt = stmt.where(DrugShipment.site_name.ilike(f"%{site_name}%"))
if tracking_no:
stmt = stmt.where(DrugShipment.tracking_no.ilike(f"%{tracking_no}%"))
if status:
stmt = stmt.where(DrugShipment.status == status)
stmt = stmt.order_by(DrugShipment.created_at.desc()).offset(skip).limit(limit)
result = await db.execute(stmt)
return result.scalars().all()
async def update_shipment(
db: AsyncSession, shipment: DrugShipment, shipment_in: DrugShipmentUpdate
) -> DrugShipment:
update_data = shipment_in.model_dump(exclude_unset=True)
for key, value in update_data.items():
setattr(shipment, key, value)
await db.commit()
await db.refresh(shipment)
return shipment
async def delete_shipment(db: AsyncSession, shipment: DrugShipment) -> None:
await db.delete(shipment)
await db.commit()
+2 -142
View File
@@ -1,150 +1,10 @@
import uuid
from datetime import date, datetime, timezone
from decimal import Decimal
from typing import Sequence
from datetime import date
from sqlalchemy import Numeric, func, select, update as sa_update, case
from sqlalchemy import case, func, select
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.finance import FinanceItem
from app.models.site import Site
from app.models.subject import Subject
from app.schemas.finance import FinanceCreate, FinanceStatusUpdate, FinanceUpdate
VALID_TRANSITIONS = {
"DRAFT": {"SUBMITTED"},
"SUBMITTED": {"APPROVED", "REJECTED"},
"APPROVED": {"PAID"},
"REJECTED": set(),
"PAID": set(),
}
async def _validate_site_subject(db: AsyncSession, study_id: uuid.UUID, site_id: uuid.UUID | None, subject_id: uuid.UUID | None):
subj_site = None
if site_id:
result = await db.execute(select(Site).where(Site.id == site_id))
site = result.scalar_one_or_none()
if not site or site.study_id != study_id:
raise ValueError("Site not found in study")
if subject_id:
result = await db.execute(select(Subject).where(Subject.id == subject_id))
subj = result.scalar_one_or_none()
if not subj or subj.study_id != study_id:
raise ValueError("Subject not found in study")
subj_site = subj.site_id
if site_id and subj_site and site_id != subj_site:
raise ValueError("Site and subject mismatch")
async def create_item(db: AsyncSession, study_id: uuid.UUID, item_in: FinanceCreate, *, created_by: uuid.UUID) -> FinanceItem:
await _validate_site_subject(db, study_id, item_in.site_id, item_in.subject_id)
item = FinanceItem(
study_id=study_id,
site_id=item_in.site_id,
subject_id=item_in.subject_id,
visit_id=item_in.visit_id,
category=item_in.category,
title=item_in.title,
description=item_in.description,
currency=item_in.currency,
amount=item_in.amount,
occur_date=item_in.occur_date,
status="DRAFT",
created_by=created_by,
)
db.add(item)
await db.commit()
await db.refresh(item)
return item
async def get_item(db: AsyncSession, item_id: uuid.UUID) -> FinanceItem | None:
result = await db.execute(select(FinanceItem).where(FinanceItem.id == item_id))
return result.scalar_one_or_none()
async def list_items(
db: AsyncSession,
study_id: uuid.UUID,
*,
status: str | None = None,
category: str | None = None,
site_id: uuid.UUID | None = None,
subject_id: uuid.UUID | None = None,
date_from: date | None = None,
date_to: date | None = None,
skip: int = 0,
limit: int = 100,
) -> Sequence[FinanceItem]:
stmt = select(FinanceItem).where(FinanceItem.study_id == study_id)
if status:
stmt = stmt.where(FinanceItem.status == status)
if category:
stmt = stmt.where(FinanceItem.category == category)
if site_id:
stmt = stmt.where(FinanceItem.site_id == site_id)
if subject_id:
stmt = stmt.where(FinanceItem.subject_id == subject_id)
if date_from:
stmt = stmt.where(FinanceItem.occur_date >= date_from)
if date_to:
stmt = stmt.where(FinanceItem.occur_date <= date_to)
stmt = stmt.offset(skip).limit(limit)
result = await db.execute(stmt)
return result.scalars().all()
async def update_item_draft(db: AsyncSession, item: FinanceItem, item_in: FinanceUpdate) -> FinanceItem:
if item.status != "DRAFT":
raise ValueError("Only DRAFT items can be edited")
update_data = item_in.model_dump(exclude_unset=True)
if update_data:
await db.execute(
sa_update(FinanceItem)
.where(FinanceItem.id == item.id)
.values(**update_data)
)
await db.commit()
await db.refresh(item)
return item
async def change_status(db: AsyncSession, item: FinanceItem, status_in: FinanceStatusUpdate, *, operator_id: uuid.UUID) -> FinanceItem:
target = status_in.status
allowed = VALID_TRANSITIONS.get(item.status, set())
if target not in allowed:
raise ValueError(f"Invalid status transition {item.status} -> {target}")
update_data = {"status": target}
now = datetime.now(timezone.utc)
if target == "SUBMITTED":
update_data["submitted_at"] = now
if target == "APPROVED":
update_data["approved_at"] = now
update_data["approver_id"] = operator_id
update_data["reject_reason"] = None
if target == "REJECTED":
if not status_in.reject_reason:
raise ValueError("reject_reason required")
update_data["rejected_at"] = now
update_data["approver_id"] = operator_id
update_data["reject_reason"] = status_in.reject_reason
if target == "PAID":
if item.status != "APPROVED":
raise ValueError("Only APPROVED can transition to PAID")
update_data["paid_at"] = now
update_data["payer_id"] = operator_id
await db.execute(
sa_update(FinanceItem)
.where(FinanceItem.id == item.id)
.values(**update_data)
)
await db.commit()
await db.refresh(item)
return item
async def summary(
db: AsyncSession,
study_id: uuid.UUID,
+69
View File
@@ -0,0 +1,69 @@
import uuid
from typing import Sequence
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.finance_contract import FinanceContract
from app.schemas.finance_contract import FinanceContractCreate, FinanceContractUpdate
async def create_contract(
db: AsyncSession,
study_id: uuid.UUID,
contract_in: FinanceContractCreate,
created_by: uuid.UUID | None,
) -> FinanceContract:
contract = FinanceContract(
study_id=study_id,
site_name=contract_in.site_name,
contract_no=contract_in.contract_no,
signed_date=contract_in.signed_date,
amount=contract_in.amount,
currency=contract_in.currency,
remark=contract_in.remark,
created_by=created_by,
)
db.add(contract)
await db.commit()
await db.refresh(contract)
return contract
async def get_contract(db: AsyncSession, contract_id: uuid.UUID) -> FinanceContract | None:
result = await db.execute(select(FinanceContract).where(FinanceContract.id == contract_id))
return result.scalar_one_or_none()
async def list_contracts(
db: AsyncSession,
study_id: uuid.UUID,
site_name: str | None = None,
contract_no: str | None = None,
skip: int = 0,
limit: int = 100,
) -> Sequence[FinanceContract]:
stmt = select(FinanceContract).where(FinanceContract.study_id == study_id)
if site_name:
stmt = stmt.where(FinanceContract.site_name.ilike(f"%{site_name}%"))
if contract_no:
stmt = stmt.where(FinanceContract.contract_no.ilike(f"%{contract_no}%"))
stmt = stmt.order_by(FinanceContract.created_at.desc()).offset(skip).limit(limit)
result = await db.execute(stmt)
return result.scalars().all()
async def update_contract(
db: AsyncSession, contract: FinanceContract, contract_in: FinanceContractUpdate
) -> FinanceContract:
update_data = contract_in.model_dump(exclude_unset=True)
for key, value in update_data.items():
setattr(contract, key, value)
await db.commit()
await db.refresh(contract)
return contract
async def delete_contract(db: AsyncSession, contract: FinanceContract) -> None:
await db.delete(contract)
await db.commit()
+69
View File
@@ -0,0 +1,69 @@
import uuid
from typing import Sequence
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.finance_special import FinanceSpecial
from app.schemas.finance_special import FinanceSpecialCreate, FinanceSpecialUpdate
async def create_special(
db: AsyncSession,
study_id: uuid.UUID,
special_in: FinanceSpecialCreate,
created_by: uuid.UUID | None,
) -> FinanceSpecial:
special = FinanceSpecial(
study_id=study_id,
site_name=special_in.site_name,
fee_type=special_in.fee_type,
amount=special_in.amount,
occur_date=special_in.occur_date,
staff_name=special_in.staff_name,
remark=special_in.remark,
created_by=created_by,
)
db.add(special)
await db.commit()
await db.refresh(special)
return special
async def get_special(db: AsyncSession, special_id: uuid.UUID) -> FinanceSpecial | None:
result = await db.execute(select(FinanceSpecial).where(FinanceSpecial.id == special_id))
return result.scalar_one_or_none()
async def list_specials(
db: AsyncSession,
study_id: uuid.UUID,
site_name: str | None = None,
fee_type: str | None = None,
skip: int = 0,
limit: int = 100,
) -> Sequence[FinanceSpecial]:
stmt = select(FinanceSpecial).where(FinanceSpecial.study_id == study_id)
if site_name:
stmt = stmt.where(FinanceSpecial.site_name.ilike(f"%{site_name}%"))
if fee_type:
stmt = stmt.where(FinanceSpecial.fee_type == fee_type)
stmt = stmt.order_by(FinanceSpecial.created_at.desc()).offset(skip).limit(limit)
result = await db.execute(stmt)
return result.scalars().all()
async def update_special(
db: AsyncSession, special: FinanceSpecial, special_in: FinanceSpecialUpdate
) -> FinanceSpecial:
update_data = special_in.model_dump(exclude_unset=True)
for key, value in update_data.items():
setattr(special, key, value)
await db.commit()
await db.refresh(special)
return special
async def delete_special(db: AsyncSession, special: FinanceSpecial) -> None:
await db.delete(special)
await db.commit()
-57
View File
@@ -1,57 +0,0 @@
import uuid
from typing import Sequence
from sqlalchemy import select, update as sa_update
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.imp_batch import ImpBatch
from app.models.imp_product import ImpProduct
from app.schemas.imp import BatchCreate, BatchUpdate
async def create_batch(db: AsyncSession, study_id: uuid.UUID, batch_in: BatchCreate) -> ImpBatch:
result = await db.execute(select(ImpProduct).where(ImpProduct.id == batch_in.product_id))
product = result.scalar_one_or_none()
if not product or product.study_id != study_id:
raise ValueError("Product not found in study")
batch = ImpBatch(
study_id=study_id,
product_id=batch_in.product_id,
batch_no=batch_in.batch_no,
expiry_date=batch_in.expiry_date,
manufacture_date=batch_in.manufacture_date,
status="ACTIVE",
)
db.add(batch)
await db.commit()
await db.refresh(batch)
return batch
async def get_batch(db: AsyncSession, batch_id: uuid.UUID) -> ImpBatch | None:
result = await db.execute(select(ImpBatch).where(ImpBatch.id == batch_id))
return result.scalar_one_or_none()
async def list_batches(db: AsyncSession, study_id: uuid.UUID, product_id: uuid.UUID | None = None, status: str | None = None) -> Sequence[ImpBatch]:
stmt = select(ImpBatch).where(ImpBatch.study_id == study_id)
if product_id:
stmt = stmt.where(ImpBatch.product_id == product_id)
if status:
stmt = stmt.where(ImpBatch.status == status)
result = await db.execute(stmt)
return result.scalars().all()
async def update_batch(db: AsyncSession, batch: ImpBatch, batch_in: BatchUpdate) -> ImpBatch:
update_data = batch_in.model_dump(exclude_unset=True)
if update_data:
await db.execute(
sa_update(ImpBatch)
.where(ImpBatch.id == batch.id)
.values(**update_data)
)
await db.commit()
await db.refresh(batch)
return batch
-46
View File
@@ -1,46 +0,0 @@
import uuid
from typing import Sequence
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.imp_inventory import ImpInventory
from app.models.imp_batch import ImpBatch
async def get_inventory_for_update(db: AsyncSession, study_id: uuid.UUID, site_id: uuid.UUID, batch_id: uuid.UUID) -> ImpInventory | None:
result = await db.execute(
select(ImpInventory)
.where(
ImpInventory.study_id == study_id,
ImpInventory.site_id == site_id,
ImpInventory.batch_id == batch_id,
)
.with_for_update()
)
return result.scalar_one_or_none()
async def get_or_create_inventory(db: AsyncSession, study_id: uuid.UUID, site_id: uuid.UUID, batch_id: uuid.UUID) -> ImpInventory:
inv = await get_inventory_for_update(db, study_id, site_id, batch_id)
if inv:
return inv
inv = ImpInventory(study_id=study_id, site_id=site_id, batch_id=batch_id, quantity_on_hand=0)
db.add(inv)
await db.flush()
return inv
async def list_inventory(
db: AsyncSession,
study_id: uuid.UUID,
site_id: uuid.UUID | None = None,
product_id: uuid.UUID | None = None,
) -> Sequence[ImpInventory]:
stmt = select(ImpInventory).where(ImpInventory.study_id == study_id)
if site_id:
stmt = stmt.where(ImpInventory.site_id == site_id)
if product_id:
stmt = stmt.join(ImpBatch, ImpBatch.id == ImpInventory.batch_id).where(ImpBatch.product_id == product_id)
result = await db.execute(stmt)
return result.scalars().all()
-50
View File
@@ -1,50 +0,0 @@
import uuid
from typing import Sequence
from sqlalchemy import select, update as sa_update
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.imp_product import ImpProduct
from app.schemas.imp import ProductCreate, ProductUpdate
async def create_product(db: AsyncSession, study_id: uuid.UUID, product_in: ProductCreate) -> ImpProduct:
product = ImpProduct(
study_id=study_id,
name=product_in.name,
form=product_in.form,
strength=product_in.strength,
unit=product_in.unit,
description=product_in.description,
is_active=product_in.is_active,
)
db.add(product)
await db.commit()
await db.refresh(product)
return product
async def get_product(db: AsyncSession, product_id: uuid.UUID) -> ImpProduct | None:
result = await db.execute(select(ImpProduct).where(ImpProduct.id == product_id))
return result.scalar_one_or_none()
async def list_products(db: AsyncSession, study_id: uuid.UUID, is_active: bool | None = None) -> Sequence[ImpProduct]:
stmt = select(ImpProduct).where(ImpProduct.study_id == study_id)
if is_active is not None:
stmt = stmt.where(ImpProduct.is_active == is_active)
result = await db.execute(stmt)
return result.scalars().all()
async def update_product(db: AsyncSession, product: ImpProduct, product_in: ProductUpdate) -> ImpProduct:
update_data = product_in.model_dump(exclude_unset=True)
if update_data:
await db.execute(
sa_update(ImpProduct)
.where(ImpProduct.id == product.id)
.values(**update_data)
)
await db.commit()
await db.refresh(product)
return product
-70
View File
@@ -1,70 +0,0 @@
import uuid
from datetime import date
from typing import Sequence
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.imp_transaction import ImpTransaction
async def create_tx(
db: AsyncSession,
*,
study_id: uuid.UUID,
site_id: uuid.UUID,
batch_id: uuid.UUID,
subject_id: uuid.UUID | None,
tx_type: str,
quantity: int,
tx_date: date,
reference: str | None,
notes: str | None,
created_by: uuid.UUID,
) -> ImpTransaction:
tx = ImpTransaction(
study_id=study_id,
site_id=site_id,
batch_id=batch_id,
subject_id=subject_id,
tx_type=tx_type,
quantity=quantity,
tx_date=tx_date,
reference=reference,
notes=notes,
created_by=created_by,
)
db.add(tx)
await db.flush()
return tx
async def list_txs(
db: AsyncSession,
study_id: uuid.UUID,
*,
site_id: uuid.UUID | None = None,
batch_id: uuid.UUID | None = None,
subject_id: uuid.UUID | None = None,
tx_type: str | None = None,
date_from: date | None = None,
date_to: date | None = None,
skip: int = 0,
limit: int = 100,
) -> Sequence[ImpTransaction]:
stmt = select(ImpTransaction).where(ImpTransaction.study_id == study_id)
if site_id:
stmt = stmt.where(ImpTransaction.site_id == site_id)
if batch_id:
stmt = stmt.where(ImpTransaction.batch_id == batch_id)
if subject_id:
stmt = stmt.where(ImpTransaction.subject_id == subject_id)
if tx_type:
stmt = stmt.where(ImpTransaction.tx_type == tx_type)
if date_from:
stmt = stmt.where(ImpTransaction.tx_date >= date_from)
if date_to:
stmt = stmt.where(ImpTransaction.tx_date <= date_to)
stmt = stmt.offset(skip).limit(limit)
result = await db.execute(stmt)
return result.scalars().all()
-99
View File
@@ -1,99 +0,0 @@
import uuid
from datetime import date, datetime, timezone
from typing import Sequence
from sqlalchemy import select, update as sa_update
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.issue import Issue
from app.models.site import Site
from app.models.subject import Subject
from app.schemas.issue import IssueCreate, IssueUpdate
async def _validate_site_subject(db: AsyncSession, study_id: uuid.UUID, site_id: uuid.UUID | None, subject_id: uuid.UUID | None):
if site_id:
result = await db.execute(select(Site).where(Site.id == site_id))
site = result.scalar_one_or_none()
if not site or site.study_id != study_id:
raise ValueError("Site not found in study")
if subject_id:
result = await db.execute(select(Subject).where(Subject.id == subject_id))
subj = result.scalar_one_or_none()
if not subj or subj.study_id != study_id:
raise ValueError("Subject not found in study")
async def create_issue(
db: AsyncSession,
study_id: uuid.UUID,
issue_in: IssueCreate,
*,
created_by: uuid.UUID,
) -> Issue:
await _validate_site_subject(db, study_id, issue_in.site_id, issue_in.subject_id)
issue = Issue(
study_id=study_id,
site_id=issue_in.site_id,
subject_id=issue_in.subject_id,
title=issue_in.title,
description=issue_in.description,
category=issue_in.category,
level=issue_in.level,
owner_id=issue_in.owner_id,
due_date=issue_in.due_date,
status="OPEN",
capa=None,
closed_at=None,
created_by=created_by,
)
db.add(issue)
await db.commit()
await db.refresh(issue)
return issue
async def get_issue(db: AsyncSession, issue_id: uuid.UUID) -> Issue | None:
result = await db.execute(select(Issue).where(Issue.id == issue_id))
return result.scalar_one_or_none()
async def list_issues(
db: AsyncSession,
study_id: uuid.UUID,
status: str | None = None,
level: str | None = None,
category: str | None = None,
overdue: bool | None = None,
) -> Sequence[Issue]:
stmt = select(Issue).where(Issue.study_id == study_id)
if status:
stmt = stmt.where(Issue.status == status)
if level:
stmt = stmt.where(Issue.level == level)
if category:
stmt = stmt.where(Issue.category == category)
if overdue is True:
stmt = stmt.where(Issue.due_date < date.today(), Issue.status != "CLOSED")
if overdue is False:
stmt = stmt.where((Issue.due_date >= date.today()) | (Issue.due_date.is_(None)) | (Issue.status == "CLOSED"))
result = await db.execute(stmt)
return result.scalars().all()
async def update_issue(db: AsyncSession, issue: Issue, issue_in: IssueUpdate) -> Issue:
update_data = issue_in.model_dump(exclude_unset=True)
if "status" in update_data:
if update_data["status"] == "CLOSED":
update_data["closed_at"] = datetime.now(timezone.utc)
else:
update_data["closed_at"] = None
if update_data:
await db.execute(
sa_update(Issue)
.where(Issue.id == issue.id)
.values(**update_data)
)
await db.commit()
await db.refresh(issue)
return issue
+67
View File
@@ -0,0 +1,67 @@
import uuid
from typing import Sequence
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.knowledge_note import KnowledgeNote
from app.schemas.knowledge_note import KnowledgeNoteCreate, KnowledgeNoteUpdate
async def create_note(
db: AsyncSession,
study_id: uuid.UUID,
note_in: KnowledgeNoteCreate,
created_by: uuid.UUID | None,
) -> KnowledgeNote:
note = KnowledgeNote(
study_id=study_id,
site_name=note_in.site_name,
title=note_in.title,
content=note_in.content,
level=note_in.level,
created_by=created_by,
)
db.add(note)
await db.commit()
await db.refresh(note)
return note
async def get_note(db: AsyncSession, note_id: uuid.UUID) -> KnowledgeNote | None:
result = await db.execute(select(KnowledgeNote).where(KnowledgeNote.id == note_id))
return result.scalar_one_or_none()
async def list_notes(
db: AsyncSession,
study_id: uuid.UUID,
site_name: str | None = None,
keyword: str | None = None,
skip: int = 0,
limit: int = 100,
) -> Sequence[KnowledgeNote]:
stmt = select(KnowledgeNote).where(KnowledgeNote.study_id == study_id)
if site_name:
stmt = stmt.where(KnowledgeNote.site_name.ilike(f"%{site_name}%"))
if keyword:
stmt = stmt.where(KnowledgeNote.title.ilike(f"%{keyword}%"))
stmt = stmt.order_by(KnowledgeNote.updated_at.desc()).offset(skip).limit(limit)
result = await db.execute(stmt)
return result.scalars().all()
async def update_note(
db: AsyncSession, note: KnowledgeNote, note_in: KnowledgeNoteUpdate
) -> KnowledgeNote:
update_data = note_in.model_dump(exclude_unset=True)
for key, value in update_data.items():
setattr(note, key, value)
await db.commit()
await db.refresh(note)
return note
async def delete_note(db: AsyncSession, note: KnowledgeNote) -> None:
await db.delete(note)
await db.commit()
-54
View File
@@ -1,54 +0,0 @@
import uuid
from typing import Sequence
from sqlalchemy import delete, select, update as sa_update
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.milestone import Milestone
from app.schemas.milestone import MilestoneCreate, MilestoneUpdate
async def create(db: AsyncSession, study_id: uuid.UUID, milestone_in: MilestoneCreate) -> Milestone:
milestone = Milestone(
study_id=study_id,
type=milestone_in.type,
name=milestone_in.name or milestone_in.type,
planned_date=milestone_in.planned_date,
actual_date=None,
status=milestone_in.status or "NOT_STARTED",
owner_id=milestone_in.owner_id,
site_id=milestone_in.site_id,
notes=milestone_in.notes,
)
db.add(milestone)
await db.commit()
await db.refresh(milestone)
return milestone
async def get(db: AsyncSession, milestone_id: uuid.UUID) -> Milestone | None:
result = await db.execute(select(Milestone).where(Milestone.id == milestone_id))
return result.scalar_one_or_none()
async def list_milestones(db: AsyncSession, study_id: uuid.UUID) -> Sequence[Milestone]:
result = await db.execute(select(Milestone).where(Milestone.study_id == study_id).order_by(Milestone.planned_date))
return result.scalars().all()
async def update(db: AsyncSession, milestone: Milestone, milestone_in: MilestoneUpdate) -> Milestone:
update_data = milestone_in.model_dump(exclude_unset=True)
if update_data:
await db.execute(
sa_update(Milestone)
.where(Milestone.id == milestone.id)
.values(**update_data)
)
await db.commit()
await db.refresh(milestone)
return milestone
async def delete_milestone(db: AsyncSession, milestone: Milestone) -> None:
await db.execute(delete(Milestone).where(Milestone.id == milestone.id))
await db.commit()
+257
View File
@@ -0,0 +1,257 @@
import uuid
from typing import Sequence
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.startup_feasibility import StartupFeasibility
from app.models.startup_ethics import StartupEthics
from app.models.kickoff_meeting import KickoffMeeting
from app.models.training_authorization import TrainingAuthorization
from app.schemas.startup import (
KickoffMeetingCreate,
KickoffMeetingUpdate,
StartupEthicsCreate,
StartupEthicsUpdate,
StartupFeasibilityCreate,
StartupFeasibilityUpdate,
TrainingAuthorizationCreate,
TrainingAuthorizationUpdate,
)
async def create_feasibility(
db: AsyncSession,
study_id: uuid.UUID,
record_in: StartupFeasibilityCreate,
created_by: uuid.UUID | None,
) -> StartupFeasibility:
record = StartupFeasibility(
study_id=study_id,
submit_date=record_in.submit_date,
accept_date=record_in.accept_date,
approved_date=record_in.approved_date,
project_no=record_in.project_no,
created_by=created_by,
)
db.add(record)
await db.commit()
await db.refresh(record)
return record
async def get_feasibility(db: AsyncSession, record_id: uuid.UUID) -> StartupFeasibility | None:
result = await db.execute(select(StartupFeasibility).where(StartupFeasibility.id == record_id))
return result.scalar_one_or_none()
async def list_feasibilities(
db: AsyncSession,
study_id: uuid.UUID,
skip: int = 0,
limit: int = 100,
) -> Sequence[StartupFeasibility]:
stmt = (
select(StartupFeasibility)
.where(StartupFeasibility.study_id == study_id)
.order_by(StartupFeasibility.created_at.desc())
.offset(skip)
.limit(limit)
)
result = await db.execute(stmt)
return result.scalars().all()
async def update_feasibility(
db: AsyncSession, record: StartupFeasibility, record_in: StartupFeasibilityUpdate
) -> StartupFeasibility:
update_data = record_in.model_dump(exclude_unset=True)
for key, value in update_data.items():
setattr(record, key, value)
await db.commit()
await db.refresh(record)
return record
async def delete_feasibility(db: AsyncSession, record: StartupFeasibility) -> None:
await db.delete(record)
await db.commit()
async def create_ethics(
db: AsyncSession,
study_id: uuid.UUID,
record_in: StartupEthicsCreate,
created_by: uuid.UUID | None,
) -> StartupEthics:
record = StartupEthics(
study_id=study_id,
submit_date=record_in.submit_date,
accept_date=record_in.accept_date,
meeting_date=record_in.meeting_date,
approved_date=record_in.approved_date,
approval_no=record_in.approval_no,
created_by=created_by,
)
db.add(record)
await db.commit()
await db.refresh(record)
return record
async def get_ethics(db: AsyncSession, record_id: uuid.UUID) -> StartupEthics | None:
result = await db.execute(select(StartupEthics).where(StartupEthics.id == record_id))
return result.scalar_one_or_none()
async def list_ethics(
db: AsyncSession,
study_id: uuid.UUID,
skip: int = 0,
limit: int = 100,
) -> Sequence[StartupEthics]:
stmt = (
select(StartupEthics)
.where(StartupEthics.study_id == study_id)
.order_by(StartupEthics.created_at.desc())
.offset(skip)
.limit(limit)
)
result = await db.execute(stmt)
return result.scalars().all()
async def update_ethics(
db: AsyncSession, record: StartupEthics, record_in: StartupEthicsUpdate
) -> StartupEthics:
update_data = record_in.model_dump(exclude_unset=True)
for key, value in update_data.items():
setattr(record, key, value)
await db.commit()
await db.refresh(record)
return record
async def delete_ethics(db: AsyncSession, record: StartupEthics) -> None:
await db.delete(record)
await db.commit()
async def create_kickoff(
db: AsyncSession,
study_id: uuid.UUID,
meeting_in: KickoffMeetingCreate,
created_by: uuid.UUID | None,
) -> KickoffMeeting:
meeting = KickoffMeeting(
study_id=study_id,
kickoff_date=meeting_in.kickoff_date,
attendees=meeting_in.attendees,
created_by=created_by,
)
db.add(meeting)
await db.commit()
await db.refresh(meeting)
return meeting
async def get_kickoff(db: AsyncSession, meeting_id: uuid.UUID) -> KickoffMeeting | None:
result = await db.execute(select(KickoffMeeting).where(KickoffMeeting.id == meeting_id))
return result.scalar_one_or_none()
async def list_kickoffs(
db: AsyncSession,
study_id: uuid.UUID,
skip: int = 0,
limit: int = 100,
) -> Sequence[KickoffMeeting]:
stmt = (
select(KickoffMeeting)
.where(KickoffMeeting.study_id == study_id)
.order_by(KickoffMeeting.created_at.desc())
.offset(skip)
.limit(limit)
)
result = await db.execute(stmt)
return result.scalars().all()
async def update_kickoff(
db: AsyncSession, meeting: KickoffMeeting, meeting_in: KickoffMeetingUpdate
) -> KickoffMeeting:
update_data = meeting_in.model_dump(exclude_unset=True)
for key, value in update_data.items():
setattr(meeting, key, value)
await db.commit()
await db.refresh(meeting)
return meeting
async def delete_kickoff(db: AsyncSession, meeting: KickoffMeeting) -> None:
await db.delete(meeting)
await db.commit()
async def create_training_authorization(
db: AsyncSession,
study_id: uuid.UUID,
record_in: TrainingAuthorizationCreate,
created_by: uuid.UUID | None,
) -> TrainingAuthorization:
record = TrainingAuthorization(
study_id=study_id,
name=record_in.name,
role=record_in.role,
site_name=record_in.site_name,
trained=record_in.trained,
authorized=record_in.authorized,
trained_date=record_in.trained_date,
authorized_date=record_in.authorized_date,
remark=record_in.remark,
created_by=created_by,
)
db.add(record)
await db.commit()
await db.refresh(record)
return record
async def get_training_authorization(
db: AsyncSession, record_id: uuid.UUID
) -> TrainingAuthorization | None:
result = await db.execute(select(TrainingAuthorization).where(TrainingAuthorization.id == record_id))
return result.scalar_one_or_none()
async def list_training_authorizations(
db: AsyncSession,
study_id: uuid.UUID,
skip: int = 0,
limit: int = 200,
) -> Sequence[TrainingAuthorization]:
stmt = (
select(TrainingAuthorization)
.where(TrainingAuthorization.study_id == study_id)
.order_by(TrainingAuthorization.created_at.desc())
.offset(skip)
.limit(limit)
)
result = await db.execute(stmt)
return result.scalars().all()
async def update_training_authorization(
db: AsyncSession, record: TrainingAuthorization, record_in: TrainingAuthorizationUpdate
) -> TrainingAuthorization:
update_data = record_in.model_dump(exclude_unset=True)
for key, value in update_data.items():
setattr(record, key, value)
await db.commit()
await db.refresh(record)
return record
async def delete_training_authorization(db: AsyncSession, record: TrainingAuthorization) -> None:
await db.delete(record)
await db.commit()
+8
View File
@@ -57,12 +57,15 @@ async def list_subjects(
study_id: uuid.UUID,
site_id: uuid.UUID | None = None,
status: str | None = None,
subject_no: str | None = None,
) -> Sequence[Subject]:
stmt = select(Subject).where(Subject.study_id == study_id)
if site_id:
stmt = stmt.where(Subject.site_id == site_id)
if status:
stmt = stmt.where(Subject.status == status)
if subject_no:
stmt = stmt.where(Subject.subject_no.ilike(f"%{subject_no}%"))
result = await db.execute(stmt)
return result.scalars().all()
@@ -102,3 +105,8 @@ async def update_subject(db: AsyncSession, subject: Subject, subject_in: Subject
await db.commit()
await db.refresh(subject)
return subject
async def delete_subject(db: AsyncSession, subject: Subject) -> None:
await db.delete(subject)
await db.commit()
+75
View File
@@ -0,0 +1,75 @@
import uuid
from typing import Sequence
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.subject import Subject
from app.models.subject_history import SubjectHistory
from app.schemas.subject_history import SubjectHistoryCreate, SubjectHistoryUpdate
async def _ensure_subject(db: AsyncSession, study_id: uuid.UUID, subject_id: uuid.UUID) -> None:
result = await db.execute(select(Subject).where(Subject.id == subject_id))
subject = result.scalar_one_or_none()
if not subject or subject.study_id != study_id:
raise ValueError("Subject not found in study")
async def create_history(
db: AsyncSession,
study_id: uuid.UUID,
history_in: SubjectHistoryCreate,
created_by: uuid.UUID | None,
) -> SubjectHistory:
await _ensure_subject(db, study_id, history_in.subject_id)
history = SubjectHistory(
study_id=study_id,
subject_id=history_in.subject_id,
record_date=history_in.record_date,
content=history_in.content,
created_by=created_by,
)
db.add(history)
await db.commit()
await db.refresh(history)
return history
async def get_history(db: AsyncSession, history_id: uuid.UUID) -> SubjectHistory | None:
result = await db.execute(select(SubjectHistory).where(SubjectHistory.id == history_id))
return result.scalar_one_or_none()
async def list_histories(
db: AsyncSession,
study_id: uuid.UUID,
subject_id: uuid.UUID,
skip: int = 0,
limit: int = 200,
) -> Sequence[SubjectHistory]:
stmt = (
select(SubjectHistory)
.where(SubjectHistory.study_id == study_id, SubjectHistory.subject_id == subject_id)
.order_by(SubjectHistory.record_date.desc().nullslast(), SubjectHistory.created_at.desc())
.offset(skip)
.limit(limit)
)
result = await db.execute(stmt)
return result.scalars().all()
async def update_history(
db: AsyncSession, history: SubjectHistory, history_in: SubjectHistoryUpdate
) -> SubjectHistory:
update_data = history_in.model_dump(exclude_unset=True)
for key, value in update_data.items():
setattr(history, key, value)
await db.commit()
await db.refresh(history)
return history
async def delete_history(db: AsyncSession, history: SubjectHistory) -> None:
await db.delete(history)
await db.commit()
-108
View File
@@ -1,108 +0,0 @@
import uuid
from datetime import date
from typing import Sequence
from sqlalchemy import select, update as sa_update
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.subject import Subject
from app.models.verification import VerificationProgress
from app.schemas.verification import VerificationCreate, VerificationUpdate
async def _validate_site_subject(db: AsyncSession, study_id: uuid.UUID, site_id: uuid.UUID, subject_id: uuid.UUID):
result = await db.execute(select(Subject).where(Subject.id == subject_id))
subj = result.scalar_one_or_none()
if not subj or subj.study_id != study_id:
raise ValueError("Subject not found in study")
if subj.site_id != site_id:
raise ValueError("Site does not match subject")
async def upsert_progress(
db: AsyncSession,
study_id: uuid.UUID,
data: VerificationCreate,
) -> VerificationProgress:
await _validate_site_subject(db, study_id, data.site_id, data.subject_id)
result = await db.execute(
select(VerificationProgress).where(
VerificationProgress.study_id == study_id,
VerificationProgress.subject_id == data.subject_id,
VerificationProgress.level == data.level,
)
)
existing = result.scalar_one_or_none()
if existing:
update_data = data.model_dump()
await db.execute(
sa_update(VerificationProgress)
.where(VerificationProgress.id == existing.id)
.values(**update_data)
)
await db.commit()
await db.refresh(existing)
return existing
vp = VerificationProgress(
study_id=study_id,
site_id=data.site_id,
subject_id=data.subject_id,
level=data.level,
percent=data.percent,
last_verified_at=data.last_verified_at,
verifier_id=data.verifier_id,
notes=data.notes,
)
db.add(vp)
await db.commit()
await db.refresh(vp)
return vp
async def get_progress(db: AsyncSession, study_id: uuid.UUID, subject_id: uuid.UUID, level: str) -> VerificationProgress | None:
result = await db.execute(
select(VerificationProgress).where(
VerificationProgress.study_id == study_id,
VerificationProgress.subject_id == subject_id,
VerificationProgress.level == level,
)
)
return result.scalar_one_or_none()
async def get_by_id(db: AsyncSession, verification_id: uuid.UUID) -> VerificationProgress | None:
result = await db.execute(select(VerificationProgress).where(VerificationProgress.id == verification_id))
return result.scalar_one_or_none()
async def list_progress(
db: AsyncSession,
study_id: uuid.UUID,
site_id: uuid.UUID | None = None,
subject_id: uuid.UUID | None = None,
level: str | None = None,
) -> Sequence[VerificationProgress]:
stmt = select(VerificationProgress).where(VerificationProgress.study_id == study_id)
if site_id:
stmt = stmt.where(VerificationProgress.site_id == site_id)
if subject_id:
stmt = stmt.where(VerificationProgress.subject_id == subject_id)
if level:
stmt = stmt.where(VerificationProgress.level == level)
result = await db.execute(stmt)
return result.scalars().all()
async def update_progress(db: AsyncSession, vp: VerificationProgress, data: VerificationUpdate) -> VerificationProgress:
update_data = data.model_dump(exclude_unset=True)
if update_data:
await db.execute(
sa_update(VerificationProgress)
.where(VerificationProgress.id == vp.id)
.values(**update_data)
)
await db.commit()
await db.refresh(vp)
return vp
+5
View File
@@ -58,3 +58,8 @@ async def update_visit(db: AsyncSession, visit: Visit, visit_in: VisitUpdate) ->
await db.commit()
await db.refresh(visit)
return visit
async def delete_visit(db: AsyncSession, visit: Visit) -> None:
await db.delete(visit)
await db.commit()