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: 3 additions & 0 deletions Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -17,5 +17,8 @@ ADD . /app
RUN --mount=type=cache,target=/root/.cache/uv \
uv sync --frozen

# set the environment variable for uvicorn to run the FastAPI app
ENV UV_PROJECT_ENVIRONMENT=.uvenv

# Run with uvicorn
CMD ["uv", "run", "uvicorn", "api.main:app", "--host", "0.0.0.0", "--port", "8000"]
56 changes: 56 additions & 0 deletions alembic_osm/versions/37c12e8301ee_workspace_jobs.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
"""workspace jobs

Revision ID: 37c12e8301ee
Revises: a92361f527ef
Create Date: 2026-08-12 10:03:15.757771

"""

from typing import Sequence, Union

import sqlalchemy as sa
from alembic import op
from sqlalchemy.dialects import postgresql

# revision identifiers, used by Alembic.
revision: str = "37c12e8301ee"
down_revision: Union[str, None] = "a92361f527ef"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None


def upgrade() -> None:
# ### commands auto generated by Alembic - please adjust! ###
op.create_table(
"jobs",
sa.Column("id", sa.Integer(), nullable=False),
sa.Column("job_type", sa.String(), nullable=False),
sa.Column("status", sa.String(), nullable=False),
sa.Column("request", postgresql.JSONB(astext_type=sa.Text()), nullable=False),
sa.Column(
"created_at",
sa.DateTime(),
server_default=sa.text("now()"),
nullable=False,
),
sa.Column(
"updated_at",
sa.DateTime(),
server_default=sa.text("now()"),
nullable=False,
),
sa.Column("current_task", sa.String(), nullable=True),
sa.Column("current_task_status", sa.String(), nullable=True),
sa.Column("response", postgresql.JSONB(astext_type=sa.Text()), nullable=True),
sa.Column("workspace_id", sa.Integer(), nullable=True),
sa.PrimaryKeyConstraint("id"),
)
# ### end Alembic commands ###
pass


def downgrade() -> None:
# ### commands auto generated by Alembic - please adjust! ###
op.drop_table("jobs")
# ### end Alembic commands ###
pass
31 changes: 31 additions & 0 deletions alembic_task/versions/aa17ec83af0d_workspace_import_status.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
"""workspace import status

Revision ID: aa17ec83af0d
Revises: d4e8f1a92b56
Create Date: 2026-08-13 04:58:59.982773

"""

from typing import Sequence, Union

import sqlalchemy as sa
from alembic import op
from sqlalchemy.dialects import postgresql

# revision identifiers, used by Alembic.
revision: str = "aa17ec83af0d"
down_revision: Union[str, None] = "d4e8f1a92b56"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None


def upgrade() -> None:
# ### commands auto generated by Alembic - please adjust! ###
op.add_column("workspaces", sa.Column("importStatus", sa.String(), nullable=True))
# ### end Alembic commands ###


def downgrade() -> None:
# ### commands auto generated by Alembic - please adjust! ###
op.drop_column("workspaces", "importStatus")
# ### end Alembic commands ###
4 changes: 4 additions & 0 deletions api/core/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,10 @@ class Settings(BaseSettings):

SENTRY_DSN: str = ""

# Azure Service Bus connection string and topic name for sending messages
SERVICE_BUS_CONNECTION_STRING: str = ""
SERVICE_BUS_TOPIC_NAME: str = ""

@property
def cors_origins_list(self) -> list[str]:
"""Allowed CORS origins as a list.
Expand Down
18 changes: 18 additions & 0 deletions api/core/messenger.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
import json

from azure.servicebus import ServiceBusClient, ServiceBusMessage

from api.core.config import settings


class Messenger:
def __init__(self):
self.connection_string = settings.SERVICE_BUS_CONNECTION_STRING
self.topic_name = settings.SERVICE_BUS_TOPIC_NAME

def send_message(self, message: dict):
with ServiceBusClient.from_connection_string(self.connection_string) as client:
sender = client.get_topic_sender(topic_name=self.topic_name)
with sender:
service_bus_message = ServiceBusMessage(json.dumps(message))
sender.send_messages(service_bus_message)
2 changes: 2 additions & 0 deletions api/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
from api.src.tasking.tasks.routes import router as tasking_tasks_router
from api.src.teams.routes import router as teams_router
from api.src.users.routes import router as users_router
from api.src.workspaces.jobs.routes import router as jobs_router
from api.src.workspaces.repository import WorkspaceRepository
from api.src.workspaces.routes import router as workspaces_router
from api.utils.migrations import run_migrations
Expand Down Expand Up @@ -108,6 +109,7 @@ async def lifespan(_app: FastAPI):
app.include_router(osm_router, prefix="/api/v1")
app.include_router(teams_router, prefix="/api/v1")
app.include_router(users_router, prefix="/api/v1")
app.include_router(jobs_router, prefix="/api/v1")
app.include_router(workspaces_router, prefix="/api/v1")
app.include_router(tasking_projects_router, prefix="/api/v1")
app.include_router(tasking_me_router, prefix="/api/v1")
Expand Down
137 changes: 137 additions & 0 deletions api/src/workspaces/jobs/repository.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,137 @@
from sqlalchemy import delete, select, update
from sqlalchemy.exc import IntegrityError
from sqlmodel.ext.asyncio.session import AsyncSession

from api.core.exceptions import (
AlreadyExistsException,
ForbiddenException,
NotFoundException,
)
from api.core.security import UserInfo
from api.src.workspaces.jobs.schemas import Job, JobCreate, JobPatch


class JobRepository:

def __init__(self, session: AsyncSession):
self.session = session

@staticmethod
def _accessible_workspace_ids(current_user: UserInfo) -> list[int]:
workspace_ids: list[int] = []
for ids in current_user.accessibleWorkspaceIds.values():
workspace_ids.extend(ids)
return workspace_ids

async def create(self, current_user: UserInfo, job_data: JobCreate) -> Job:
# if not current_user.isWorkspaceContributor(job_data.workspace_id):
# raise ForbiddenException(
# "User does not have permissions to create a job in that workspace."
# )
# Not needed as this may be during the initial creation. There is no API call to create a job.
# The job is created internally when a workspace is created. So, we don't need to check for permissions here.

job = Job(**job_data.model_dump())

try:
self.session.add(job)
await self.session.commit()
await self.session.refresh(job)
return job
except IntegrityError:
await self.session.rollback()
raise AlreadyExistsException(f"Job with ID {job.id} already exists")

async def getById(self, current_user: UserInfo, job_id: int) -> Job:
accessible_workspace_ids = self._accessible_workspace_ids(current_user)

query = select(Job).where(
(Job.id == job_id)
& (Job.workspace_id.in_(accessible_workspace_ids)) # type: ignore[attr-defined]
)
result = await self.session.execute(query)
job = result.scalar_one_or_none()

if not job:
raise NotFoundException(f"Job with id {job_id} not found")

return job

async def update(
self,
current_user: UserInfo,
job_id: int,
job_data: JobPatch,
ignore_permissions: bool = False,
) -> Job:
if ignore_permissions:
query = (
update(Job)
.where(Job.id == job_id) # pyright: ignore[reportArgumentType]
.values(**job_data.model_dump(exclude_unset=True))
)
result = await self.session.execute(query)
if result.rowcount != 1: # type: ignore[attr-defined]
raise NotFoundException(f"Update failed for job id {job_id}")
await self.session.commit()
return await self._getJobById(job_id)
else:
accessible_workspace_ids = self._accessible_workspace_ids(current_user)

query = (
update(Job)
.where(
(Job.id == job_id)
& (Job.workspace_id.in_(accessible_workspace_ids)) # type: ignore[attr-defined]
)
.values(**job_data.model_dump(exclude_unset=True))
)

result = await self.session.execute(query)

if result.rowcount != 1: # type: ignore[attr-defined]
raise NotFoundException(f"Update failed for job id {job_id}")

await self.session.commit()
return await self._getJobById(job_id)

async def delete(self, current_user: UserInfo, job_id: int) -> None:
accessible_workspace_ids = self._accessible_workspace_ids(current_user)

query = delete(Job).where(
(Job.id == job_id)
& (Job.workspace_id.in_(accessible_workspace_ids)) # type: ignore[attr-defined]
)

result = await self.session.execute(query)

if result.rowcount != 1: # type: ignore[attr-defined]
raise NotFoundException(f"Job delete failed for id {job_id}")

await self.session.commit()

async def getWorkspaceJobs(
self, current_user: UserInfo, workspace_id: int
) -> list[Job]:
if not current_user.isWorkspaceContributor(workspace_id):
raise ForbiddenException(
"User does not have permissions to view jobs in that workspace."
)

query = select(Job).where(
Job.workspace_id == workspace_id # pyright: ignore[reportArgumentType]
)
result = await self.session.execute(query)
return list(result.scalars().all())

async def _getJobById(self, job_id: int) -> Job:
query = select(Job).where(
Job.id == job_id # pyright: ignore[reportArgumentType]
)
result = await self.session.execute(query)
job = result.scalar_one_or_none()

if not job:
raise NotFoundException(f"Job with id {job_id} not found")

return job
24 changes: 24 additions & 0 deletions api/src/workspaces/jobs/routes.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
from fastapi import APIRouter, Depends
from sqlmodel.ext.asyncio.session import AsyncSession

from api.core.database import get_osm_session
from api.core.security import UserInfo, validate_token
from api.src.workspaces.jobs.repository import JobRepository
from api.src.workspaces.jobs.schemas import Job

router = APIRouter(prefix="/workspaces/jobs", tags=["jobs"])


def get_job_repository(
session: AsyncSession = Depends(get_osm_session),
) -> JobRepository:
return JobRepository(session)


@router.get("/{job_id}", response_model=Job)
async def get_job_by_id(
job_id: int,
repository: JobRepository = Depends(get_job_repository),
current_user: UserInfo = Depends(validate_token),
) -> Job:
return await repository.getById(current_user, job_id)
51 changes: 51 additions & 0 deletions api/src/workspaces/jobs/schemas.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
from datetime import datetime
from typing import Any

from sqlalchemy import JSON as SAJson
from sqlalchemy import Column
from sqlmodel import Field, Relationship, SQLModel


class Job(SQLModel, table=True):
__tablename__ = "jobs" # type: ignore[assignment]

id: int = Field(default=None, primary_key=True)
job_type: str = Field(default=None, nullable=False)
status: str = Field(default=None, nullable=False)
request: dict[str, Any] = Field(
default=None, sa_column=Column(SAJson, nullable=False)
)
created_at: datetime = Field(sa_column=Column(nullable=False, default=datetime.now))
updated_at: datetime = Field(sa_column=Column(nullable=False, default=datetime.now))
current_task: str | None = Field(default=None, nullable=True)
current_task_status: str | None = Field(default=None, nullable=True)
response: dict[str, Any] | None = Field(
default=None, sa_column=Column(SAJson, nullable=True)
)
workspace_id: int = Field(default=None) # Not sure if we have to add index here.

# workspace_id: int = Field(foreign_key="workspaces.id")
# workspace: "Workspace" = Relationship(back_populates="jobs")


class JobCreate(SQLModel):
"""Fields the client may supply when creating a job."""

job_type: str
status: str
request: dict[str, Any]
workspace_id: int
current_task: str | None = None
current_task_status: str | None = None
response: dict[str, Any] | None = None


class JobPatch(SQLModel):
"""Fields the client may supply when updating a job."""

job_type: str | None = None
status: str | None = None
request: dict[str, Any] | None = None
current_task: str | None = None
current_task_status: str | None = None
response: dict[str, Any] | None = None
5 changes: 5 additions & 0 deletions api/src/workspaces/repository.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
WorkspaceImagery,
WorkspaceLongQuest,
WorkspacePatch,
WorkspaceType,
)


Expand All @@ -29,10 +30,14 @@ def __init__(self, session: AsyncSession):
async def create(
self, current_user: UserInfo, workspace_data: WorkspaceCreate
) -> Workspace:
importStatus = "NA"
if workspace_data.isTDEIDataset():
importStatus = "in-progress"
workspace = Workspace(
**workspace_data.model_dump(),
createdBy=current_user.user_uuid, # type: ignore[reportArgumentType]
createdByName=current_user.user_name,
importStatus=importStatus,
)

if str(workspace.tdeiProjectGroupId) not in current_user.getProjectGroupIds():
Expand Down
Loading
Loading