Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
3 changes: 1 addition & 2 deletions backend/druks/api/runs.py
Original file line number Diff line number Diff line change
Expand Up @@ -90,15 +90,14 @@ async def cancel_run(
raise RunNotActive(run_id)
subject = await run.get_subject()
# cancel() flushes the run and expires this computed column, so read it first.
key, title = run.subject_key, run.subject_title
key = run.subject_key
await run.cancel(failure=reason)
if subject:
await Event.emit(
session,
type=WorkflowEvent.CANCELLED,
subject=subject,
key=key,
title=title,
run=run.id,
kind=run.kind,
facts={"failure": reason},
Expand Down
25 changes: 19 additions & 6 deletions backend/druks/apps/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -417,13 +417,15 @@ def _get_page_routes(cls) -> "APIRouter":
from druks.ui import Page

operations = cls.operations()
declarations = cls.pages()
landing = {declaration.route: declaration for declaration in declarations}.get("/")
router = APIRouter(prefix="/pages", tags=[f"{cls.name}:pages"])
for declaration in cls.pages():
for declaration in declarations:
router.add_api_route(
# The landing page's route is "/", and its snapshot answers at
# the bare /pages.
declaration.route.rstrip("/"),
cls._page_endpoint(declaration, operations),
cls._page_endpoint(declaration, operations, back=declaration.parent or landing),
methods=["GET"],
response_model=Page,
response_model_by_alias=True,
Expand All @@ -432,17 +434,28 @@ def _get_page_routes(cls) -> "APIRouter":
return router

@classmethod
def _page_endpoint(cls, declaration: "PageRoute", operations: "dict[str, Operation]"):
def _page_endpoint(
cls,
declaration: "PageRoute",
operations: "dict[str, Operation]",
*,
back: "PageRoute | None",
):
"""``wraps`` keeps the page function's signature, so FastAPI still
validates every route parameter."""
from druks.ui import EmptyState, Page
validates every route parameter. A page whose subject is missing answers
an empty state that links ``back``."""
from druks.ui import EmptyState, Link, Page

controls = []
if back and back is not declaration:
controls.append(Link(back.label, page=back.name))

@wraps(declaration.function)
async def read_page(**parameters):
try:
page = await declaration.function(**parameters)
except ObjectNotFound as error:
return Page(str(error), blocks=[EmptyState(str(error))])
return Page(str(error), blocks=[EmptyState(str(error), controls=controls)])
except Exception as error:
raise PageReadError(
cls.name, declaration.name, f"its own code raised {type(error).__name__}"
Expand Down
4 changes: 0 additions & 4 deletions backend/druks/contrib/software_factory/datastructures.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,10 +9,6 @@ class PullRequest(Subject):
is the whole record — ``owner/repo#7`` is what everything else reads back out of,
and its status carries the lifecycle."""

@classmethod
def get(cls, repo: str, number: int) -> Self:
return cls(id=f"{repo}#{number}")

@classmethod
async def get_or_none(cls, id: str) -> Self | None:
# Ids reach the read side as free text off a URL, so a shape that names no
Expand Down
6 changes: 3 additions & 3 deletions backend/druks/contrib/software_factory/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -126,7 +126,7 @@ async def get_in_project(cls, *, project_id: int, repo_id: int) -> "ProjectRepo
stmt = select(cls).where(cls.id == repo_id, cls.project_id == project_id).limit(1)
return (await db_session().scalars(stmt)).first()

def get_key(self) -> str:
def __str__(self) -> str:
return self.full_name

def get_summary(self) -> "ProjectRepoSummary":
Expand Down Expand Up @@ -234,8 +234,8 @@ class WorkItem(StoredSubject):
# The time of the GitHub verdict, or of the cancel reaction.
resolved_at: Mapped[datetime | None] = mapped_column(default=None)

def get_key(self) -> str:
return self.ticket_key
def __str__(self) -> str:
return f"{self.ticket_key} {self.title}".strip()

def get_summary(self) -> WorkItemSummary:
return WorkItemSummary.model_validate(self)
Expand Down
2 changes: 1 addition & 1 deletion backend/druks/contrib/software_factory/workflows.py
Original file line number Diff line number Diff line change
Expand Up @@ -531,7 +531,7 @@ async def dispatch(cls, *, repo: str, pr_number: int, account: Account, note: st
# reviewer is connected. The lookup raises a clear error before the run starts a VM.
await Github.get()
return await cls.start(
subject=PullRequest.get(repo, pr_number), account_id=account.id, note=note
subject=PullRequest(id=f"{repo}#{pr_number}"), account_id=account.id, note=note
)

async def run(self, note: str = "") -> None:
Expand Down
28 changes: 20 additions & 8 deletions backend/druks/durable/datastructures.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
from druks.db import db_session
from druks.durable.schemas import SubjectSummary
from druks.events.models import Event
from druks.exceptions import ObjectNotFound
from druks.models import snake_name

if TYPE_CHECKING:
Expand All @@ -22,9 +23,6 @@ class Subject:
instead."""

subject_type: ClassVar[str]
# The header its board and page show it under. Set a ``SubjectSummary``
# subclass to add the app's own fields and a descriptive ``title``.
summary_class: ClassVar[type[SubjectSummary]] = SubjectSummary

id: str

Expand All @@ -38,11 +36,14 @@ def __init_subclass__(cls, **kwargs: Any) -> None:
def identity(self) -> dict[str, Any]:
return {"type": self.subject_type, "id": self.id}

def __str__(self) -> str:
"""How the subject shows itself. An identity-only subject is already named by
its id: "owner/repo#7" is the handle, not a surrogate key."""
return self.id

@property
def key(self) -> str:
# An identity-only subject is already named by its id — "owner/repo#7" is
# the handle, not a surrogate key.
return self.id
return str(self)

async def announce(self, topic: str, **facts: Any) -> None:
"""Record and deliver a domain fact in the current transaction."""
Expand All @@ -54,12 +55,23 @@ async def get_or_none(cls, id: str) -> Self | None:
so override to return None for a shape this subject could never wear."""
return cls(id=id)

@classmethod
async def get(cls, id: str) -> Self:
"""The subject this id names. A shape it could never wear raises
``ObjectNotFound``, which a route answers with 404 and a page with an empty
state."""
if subject := await cls.get_or_none(id):
return subject
raise ObjectNotFound(cls.subject_type.replace("_", " "), {"id": id})

def get_summary(self) -> SubjectSummary:
return self.summary_class.model_validate(self)
"""The header the platform's own screens show. An app with its own frontend
overrides it to add the fields that frontend reads."""
return SubjectSummary.model_validate(self)

@classmethod
async def list_summaries(cls, account_id: str | None) -> Sequence[SubjectSummary]:
"""The subjects on this class's board, newest-movement first, each as its domain
"""The subjects on this class's board, newest movement first, each as its
summary. ``account_id`` is the caller, or None outside a request. A shared
board ignores it. Returns a covariant ``Sequence`` so an app can return
a ``list`` of its own ``SubjectSummary`` subclass. Required once a workflow
Expand Down
17 changes: 11 additions & 6 deletions backend/druks/durable/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -77,9 +77,6 @@ class Run(Base):
subject_key: Mapped[str | None] = column_property(
subject_attribute_expression(id, "subject_key")
)
subject_title: Mapped[str | None] = column_property(
subject_attribute_expression(id, "subject_title")
)
retry_from: Mapped[str | None] = column_property(retry_from_expression(id))
account_id: Mapped[str] = mapped_column(ForeignKey("accounts.id", ondelete="RESTRICT"))
account: Mapped[Account] = relationship(lazy="joined", foreign_keys=[account_id])
Expand Down Expand Up @@ -183,11 +180,18 @@ async def list_for_subject(

@classmethod
async def get_latest_for_subject(
cls, session: AsyncSession, subject_type: str, subject_id: str, kind: str | None = None
cls,
session: AsyncSession,
subject_type: str,
subject_id: str,
kind: str | None = None,
*,
gate: str | None = None,
) -> "Run | None":
"""The run that speaks for the subject: a subject holds at most one active
run per kind (queue dedup) and the next starts only once the last is
terminal, so the newest is the live one whenever anything is live."""
terminal, so the newest is the live one whenever anything is live. ``gate``
narrows to the newest run parked on that gate."""
stmt = (
select(cls)
.where(subject_filter(cls.id, subject_type, subject_id))
Expand All @@ -197,6 +201,8 @@ async def get_latest_for_subject(
)
if kind:
stmt = stmt.where(cls.kind == kind)
if gate:
stmt = stmt.where(cls.input_gate == gate, cls.state == RunState.PARKED.value)
return (await session.scalars(stmt)).first()

@classmethod
Expand Down Expand Up @@ -734,7 +740,6 @@ async def record(
type=event["topic"],
subject=await run.get_subject(),
key=run.subject_key,
title=run.subject_title,
run=run.id,
kind=run.kind,
facts=facts,
Expand Down
7 changes: 3 additions & 4 deletions backend/druks/durable/schemas.py
Original file line number Diff line number Diff line change
Expand Up @@ -140,14 +140,13 @@ def from_run(


class SubjectSummary(Schema):
# The base an app's subject header subclasses; ``id`` keys the subject's
# status, timeline and detail URL, and ``from_attributes`` builds the header
# straight off the subject.
# The header the platform shows a subject under: ``id`` keys its status,
# timeline and detail URL, ``key`` is its name, and ``from_attributes`` builds
# it straight off the subject. An app with its own frontend subclasses it.
model_config = ConfigDict(from_attributes=True)

id: SubjectId
key: SubjectKey
title: str | None = None


class SubjectStatus(Schema):
Expand Down
13 changes: 3 additions & 10 deletions backend/druks/events/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,7 @@ def get_history(
until: datetime | None = None,
) -> Select[tuple["Event"]]:
"""The recorded Activity that matches these filters, as a query. Search reads the
recorded key and title literally; from is inclusive and until is exclusive."""
recorded key literally; from is inclusive and until is exclusive."""
# The durable package imports this module.
from druks.durable.enums import WorkflowEvent

Expand All @@ -90,12 +90,7 @@ def get_history(
if app:
statement = statement.where(cls.app == app)
if search and search.strip():
statement = statement.where(
or_(
cls.subject_key.icontains(search.strip(), autoescape=True),
cls.payload["title"].as_string().icontains(search.strip(), autoescape=True),
)
)
statement = statement.where(cls.subject_key.icontains(search.strip(), autoescape=True))
if topic:
statement = statement.where(cls.type == topic)
if from_at:
Expand Down Expand Up @@ -131,7 +126,6 @@ async def emit(
type: str,
subject: dict[str, Any] | None = None,
key: str | None = None,
title: str | None = None,
run: str | None = None,
kind: str | None = None,
facts: dict[str, Any] | None = None,
Expand All @@ -143,7 +137,7 @@ async def emit(

subject = subject or {}
facts = facts or {}
recorded = {"run": run, "kind": kind, "title": title}
recorded = {"run": run, "kind": kind}
if taken := recorded.keys() & facts.keys():
raise WorkflowError(f"{type} facts {sorted(taken)} belong to Druks. Rename them.")
session.add(
Expand Down Expand Up @@ -184,7 +178,6 @@ async def announce(
type=topic,
subject=subject.identity,
key=subject.key,
title=subject.get_summary().title,
facts=facts,
app=app,
)
Expand Down
2 changes: 1 addition & 1 deletion backend/druks/exceptions.py
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ def __init__(self) -> None:


class ObjectNotFound(DruksError):
"""No row holds the values a read asked for."""
"""Nothing matches the values a read asked for."""

def __init__(self, model: str, fields: dict[str, object]) -> None:
match = ", ".join(f"{name} {value}" for name, value in fields.items())
Expand Down
Loading
Loading