Skip to content

RPCHandler

hololinked.server.http.handlers.RPCHandler

Bases: BaseHandler

Handler for property read-write and method calls.

Source code in repo/hololinked/hololinked/server/http/handlers.py
class RPCHandler(BaseHandler):
    """Handler for property read-write and method calls."""

    async def is_method_allowed(self, method: str) -> bool:
        """
        Checks if the method is allowed for the property:

        - Access control (authentication & authorization)
        - if the HTTP method is allowed for the resource.
        - if its GET method with message id for no-block response.

        Returns
        -------
        bool
            `True` if the method may be served, `False` if a 401/403/405 was already set on the response
        """  # noqa: D400
        if not await self.has_access_control():
            return False
        if self.message_id is not None and method.upper() == "GET":
            return True
        if method in self.metadata.http_methods:
            return True
        self.set_status(405, "method not allowed")
        return False

    async def options(self) -> None:  # ty: ignore[invalid-method-override]
        """
        Options for the resource.

        Main functionality is to inform the client is a specific HTTP method is supported by
        the property or the action (Access-Control-Allow-Methods).
        """
        if await self.has_access_control():
            self.set_status(204)
            self.set_custom_default_headers()
            self.set_access_control_allow_headers()
            self.set_header("Access-Control-Allow-Methods", ", ".join(self.metadata.http_methods))
        self.finish()

    async def write_reply(self, reply: Reply) -> None:
        """Write the event loop's reply onto the wire."""
        # only one payload reaches the client for now, no support for multipart # TODO
        if reply.preserialized_payload.value:
            if reply.payload.value is not None:
                self.logger.warning(
                    "multipart payloads are not supported over HTTP, only the preserialized payload is written",
                    content_type=reply.payload.content_type,
                )
            reply_payload = reply.preserialized_payload
        else:
            reply_payload = reply.payload
        body = reply_payload.serialize() if isinstance(reply_payload, SerializableData) else reply_payload.value
        self.set_header("Content-Type", reply_payload.content_type or "application/json")
        if body:
            super().write(body)

    async def handle_through_thing(self, operation: str) -> None:
        """
        Handles the `Thing` operations and writes the reply to the HTTP client.

        Parameters
        ----------
        operation: str
            operation to be performed on the Thing, like `readproperty`,
            `writeproperty`, `invokeaction`, `deleteproperty`

        Raises
        ------
        RuntimeError
            If the `Thing` is not served by an event loop
        """
        if not self.thing.eventloop:  # type gaurd, not a real logic.
            raise RuntimeError("Thing is not served by an event loop")
        try:
            scheduler_execution_context, thing_execution_context, local_execution_context, additional_payload = (
                self.get_execution_parameters()
            )
            payload, preserialized_payload = self.get_request_payload()
            payload = payload if payload.value else additional_payload
        except Exception as ex:
            self.set_status(400, f"error while decoding request - {str(ex)}")
            self.logger.error(f"error while decoding request - {str(ex)}")
            return
        try:
            request = Operation(
                thing_id=self.thing_id,
                objekt=self.resource.name,
                operation=operation,
                payload=payload,
                preserialized_payload=preserialized_payload,
                scheduler_execution_context=scheduler_execution_context,
                thing_execution_context=thing_execution_context,
                id=uuid_hex(),
                sender_id=self.request.remote_ip or "",
            )
            if scheduler_execution_context.oneway:
                # no reply is wanted, so the future is dropped rather than awaited
                self.thing.eventloop.submit(request)
                self.set_status(204, "ok")
            elif local_execution_context.noblock:
                # the client collects this on a second request, quoting the message ID back to us
                message_id = uuid_hex()
                pending = self.thing.eventloop.pending_operations
                pending.add(self.config.server_id, message_id, self.thing.eventloop.submit(request))
                self.set_status(204, "ok")
                self.set_header("X-Message-ID", message_id)
            else:
                reply = await self.thing.eventloop.execute(request)
                if reply.timed_out:
                    self.set_status(408, f"{reply.kind.value.replace('_', ' ')} while executing the operation")
                    return
                self.set_status(200, "ok")
                await self.write_reply(reply)
        except ConnectionAbortedError as ex:
            self.set_status(503, f"lost connection to thing - {str(ex)}")
            # TODO handle reconnection
        except Exception as ex:
            self.logger.error(f"error while scheduling RPC call - {str(ex)}")
            self.set_status(500, f"error while scheduling RPC call - {str(ex)}")
            self.set_header("Content-Type", "application/json")
            response_payload = SerializableData(
                value=Serializers.json.dumps({"exception": format_exception_as_json(ex)}),
                content_type="application/json",
            )
            response_payload.serialize()
            self.write(response_payload.value)

    async def handle_no_block_response(self) -> None:
        """Handles the no-block response for the noblock calls."""  # noqa: DOC501
        if not self.thing.eventloop:  # type gaurd, not a real logic.
            raise RuntimeError("Thing is not served by an event loop")
        future = None  # held only while this request owns the claim, so that `finally` can give it back
        message_id = None
        try:
            message_id = self.message_id
            if message_id is None:
                raise ValueError("no message id available to wait for a no-block response")
            self.logger.info("waiting for no-block response", message_id=message_id)
            future = self.thing.eventloop.pending_operations.claim(self.config.server_id, message_id)
            invokation = default_scheduler_execution_context.invokation_timeout
            execution = default_scheduler_execution_context.execution_timeout
            # either being None means wait indefinitely, so there is no bound to compute
            bound = None if invokation is None or execution is None else invokation + execution
            # shielded - a timeout here must not cancel the operation, it should continue running for later collection
            reply: Reply = await asyncio.wait_for(asyncio.shield(asyncio.wrap_future(future)), timeout=bound)
            if reply.timed_out:
                self.set_status(408, f"{reply.kind.value.replace('_', ' ')} while executing the operation")
            else:
                self.set_status(200, "ok")
                await self.write_reply(reply)
            future = None  # answered - the caller has no reason to come back with this message ID
        except KeyError as ex:
            # if the message id is not found, it means that the response was not received in time
            self.logger.error(f"message ID not found for no-block response - {str(ex)}")
            self.set_status(404, "message id not found")
        except TimeoutError as ex:
            self.logger.error(f"timeout while waiting for no-block response - {str(ex)}")
            self.set_status(408, "timeout while waiting for response, ask later")
        except Exception as ex:
            self.logger.error(f"error while receiving no-block response - {str(ex)}")
            self.set_status(500, f"error while receiving no-block response - {str(ex)}")
            self.set_header("Content-Type", "application/json")
            response_payload = SerializableData(
                value=Serializers.json.dumps({"exception": format_exception_as_json(ex)}),
                content_type="application/json",
            )
            response_payload.serialize()
            self.write(response_payload.value)
        finally:
            if future is not None and message_id is not None:
                # the operation is still running
                self.thing.eventloop.pending_operations.add(self.config.server_id, message_id, future)

Functions

is_method_allowed async

is_method_allowed(method: str) -> bool

Checks if the method is allowed for the property:

  • Access control (authentication & authorization)
  • if the HTTP method is allowed for the resource.
  • if its GET method with message id for no-block response.

Returns:

Type Description
bool

True if the method may be served, False if a 401/403/405 was already set on the response

Source code in repo/hololinked/hololinked/server/http/handlers.py
async def is_method_allowed(self, method: str) -> bool:
    """
    Checks if the method is allowed for the property:

    - Access control (authentication & authorization)
    - if the HTTP method is allowed for the resource.
    - if its GET method with message id for no-block response.

    Returns
    -------
    bool
        `True` if the method may be served, `False` if a 401/403/405 was already set on the response
    """  # noqa: D400
    if not await self.has_access_control():
        return False
    if self.message_id is not None and method.upper() == "GET":
        return True
    if method in self.metadata.http_methods:
        return True
    self.set_status(405, "method not allowed")
    return False

options async

options() -> None

Options for the resource.

Main functionality is to inform the client is a specific HTTP method is supported by the property or the action (Access-Control-Allow-Methods).

Source code in repo/hololinked/hololinked/server/http/handlers.py
async def options(self) -> None:  # ty: ignore[invalid-method-override]
    """
    Options for the resource.

    Main functionality is to inform the client is a specific HTTP method is supported by
    the property or the action (Access-Control-Allow-Methods).
    """
    if await self.has_access_control():
        self.set_status(204)
        self.set_custom_default_headers()
        self.set_access_control_allow_headers()
        self.set_header("Access-Control-Allow-Methods", ", ".join(self.metadata.http_methods))
    self.finish()

handle_through_thing async

handle_through_thing(operation: str) -> None

Handles the Thing operations and writes the reply to the HTTP client.

Parameters:

Name Type Description Default
operation
str

operation to be performed on the Thing, like readproperty, writeproperty, invokeaction, deleteproperty

required

Raises:

Type Description
RuntimeError

If the Thing is not served by an event loop

Source code in repo/hololinked/hololinked/server/http/handlers.py
async def handle_through_thing(self, operation: str) -> None:
    """
    Handles the `Thing` operations and writes the reply to the HTTP client.

    Parameters
    ----------
    operation: str
        operation to be performed on the Thing, like `readproperty`,
        `writeproperty`, `invokeaction`, `deleteproperty`

    Raises
    ------
    RuntimeError
        If the `Thing` is not served by an event loop
    """
    if not self.thing.eventloop:  # type gaurd, not a real logic.
        raise RuntimeError("Thing is not served by an event loop")
    try:
        scheduler_execution_context, thing_execution_context, local_execution_context, additional_payload = (
            self.get_execution_parameters()
        )
        payload, preserialized_payload = self.get_request_payload()
        payload = payload if payload.value else additional_payload
    except Exception as ex:
        self.set_status(400, f"error while decoding request - {str(ex)}")
        self.logger.error(f"error while decoding request - {str(ex)}")
        return
    try:
        request = Operation(
            thing_id=self.thing_id,
            objekt=self.resource.name,
            operation=operation,
            payload=payload,
            preserialized_payload=preserialized_payload,
            scheduler_execution_context=scheduler_execution_context,
            thing_execution_context=thing_execution_context,
            id=uuid_hex(),
            sender_id=self.request.remote_ip or "",
        )
        if scheduler_execution_context.oneway:
            # no reply is wanted, so the future is dropped rather than awaited
            self.thing.eventloop.submit(request)
            self.set_status(204, "ok")
        elif local_execution_context.noblock:
            # the client collects this on a second request, quoting the message ID back to us
            message_id = uuid_hex()
            pending = self.thing.eventloop.pending_operations
            pending.add(self.config.server_id, message_id, self.thing.eventloop.submit(request))
            self.set_status(204, "ok")
            self.set_header("X-Message-ID", message_id)
        else:
            reply = await self.thing.eventloop.execute(request)
            if reply.timed_out:
                self.set_status(408, f"{reply.kind.value.replace('_', ' ')} while executing the operation")
                return
            self.set_status(200, "ok")
            await self.write_reply(reply)
    except ConnectionAbortedError as ex:
        self.set_status(503, f"lost connection to thing - {str(ex)}")
        # TODO handle reconnection
    except Exception as ex:
        self.logger.error(f"error while scheduling RPC call - {str(ex)}")
        self.set_status(500, f"error while scheduling RPC call - {str(ex)}")
        self.set_header("Content-Type", "application/json")
        response_payload = SerializableData(
            value=Serializers.json.dumps({"exception": format_exception_as_json(ex)}),
            content_type="application/json",
        )
        response_payload.serialize()
        self.write(response_payload.value)

write_reply async

write_reply(reply: Reply) -> None

Write the event loop's reply onto the wire.

Source code in repo/hololinked/hololinked/server/http/handlers.py
async def write_reply(self, reply: Reply) -> None:
    """Write the event loop's reply onto the wire."""
    # only one payload reaches the client for now, no support for multipart # TODO
    if reply.preserialized_payload.value:
        if reply.payload.value is not None:
            self.logger.warning(
                "multipart payloads are not supported over HTTP, only the preserialized payload is written",
                content_type=reply.payload.content_type,
            )
        reply_payload = reply.preserialized_payload
    else:
        reply_payload = reply.payload
    body = reply_payload.serialize() if isinstance(reply_payload, SerializableData) else reply_payload.value
    self.set_header("Content-Type", reply_payload.content_type or "application/json")
    if body:
        super().write(body)

handle_no_block_response async

handle_no_block_response() -> None

Handles the no-block response for the noblock calls.

Source code in repo/hololinked/hololinked/server/http/handlers.py
async def handle_no_block_response(self) -> None:
    """Handles the no-block response for the noblock calls."""  # noqa: DOC501
    if not self.thing.eventloop:  # type gaurd, not a real logic.
        raise RuntimeError("Thing is not served by an event loop")
    future = None  # held only while this request owns the claim, so that `finally` can give it back
    message_id = None
    try:
        message_id = self.message_id
        if message_id is None:
            raise ValueError("no message id available to wait for a no-block response")
        self.logger.info("waiting for no-block response", message_id=message_id)
        future = self.thing.eventloop.pending_operations.claim(self.config.server_id, message_id)
        invokation = default_scheduler_execution_context.invokation_timeout
        execution = default_scheduler_execution_context.execution_timeout
        # either being None means wait indefinitely, so there is no bound to compute
        bound = None if invokation is None or execution is None else invokation + execution
        # shielded - a timeout here must not cancel the operation, it should continue running for later collection
        reply: Reply = await asyncio.wait_for(asyncio.shield(asyncio.wrap_future(future)), timeout=bound)
        if reply.timed_out:
            self.set_status(408, f"{reply.kind.value.replace('_', ' ')} while executing the operation")
        else:
            self.set_status(200, "ok")
            await self.write_reply(reply)
        future = None  # answered - the caller has no reason to come back with this message ID
    except KeyError as ex:
        # if the message id is not found, it means that the response was not received in time
        self.logger.error(f"message ID not found for no-block response - {str(ex)}")
        self.set_status(404, "message id not found")
    except TimeoutError as ex:
        self.logger.error(f"timeout while waiting for no-block response - {str(ex)}")
        self.set_status(408, "timeout while waiting for response, ask later")
    except Exception as ex:
        self.logger.error(f"error while receiving no-block response - {str(ex)}")
        self.set_status(500, f"error while receiving no-block response - {str(ex)}")
        self.set_header("Content-Type", "application/json")
        response_payload = SerializableData(
            value=Serializers.json.dumps({"exception": format_exception_as_json(ex)}),
            content_type="application/json",
        )
        response_payload.serialize()
        self.write(response_payload.value)
    finally:
        if future is not None and message_id is not None:
            # the operation is still running
            self.thing.eventloop.pending_operations.add(self.config.server_id, message_id, future)

PropertyHandler

Bases: RPCHandler

handles property requests.

Source code in repo/hololinked/hololinked/server/http/handlers.py
class PropertyHandler(RPCHandler):
    """handles property requests."""

    async def get(self) -> None:
        """Read the property, or fetch the reply of an earlier no-block read."""
        if await self.is_method_allowed("GET"):
            self.set_custom_default_headers()
            if self.message_id is not None:
                await self.handle_no_block_response()
            else:
                await self.handle_through_thing(Operations.readproperty)
        self.finish()

    async def post(self) -> None:
        """Write the property."""
        if await self.is_method_allowed("POST"):
            self.set_custom_default_headers()
            await self.handle_through_thing(Operations.writeproperty)
        self.finish()

    async def put(self) -> None:
        """Write the property."""
        if await self.is_method_allowed("PUT"):
            self.set_custom_default_headers()
            await self.handle_through_thing(Operations.writeproperty)
        self.finish()

    async def delete(self) -> None:
        """Delete the property."""
        if await self.is_method_allowed("DELETE"):
            self.set_custom_default_headers()
            await self.handle_through_thing(Operations.deleteproperty)
        self.finish()

Functions

get async

get() -> None

Read the property, or fetch the reply of an earlier no-block read.

Source code in repo/hololinked/hololinked/server/http/handlers.py
async def get(self) -> None:
    """Read the property, or fetch the reply of an earlier no-block read."""
    if await self.is_method_allowed("GET"):
        self.set_custom_default_headers()
        if self.message_id is not None:
            await self.handle_no_block_response()
        else:
            await self.handle_through_thing(Operations.readproperty)
    self.finish()

post async

post() -> None

Write the property.

Source code in repo/hololinked/hololinked/server/http/handlers.py
async def post(self) -> None:
    """Write the property."""
    if await self.is_method_allowed("POST"):
        self.set_custom_default_headers()
        await self.handle_through_thing(Operations.writeproperty)
    self.finish()

put async

put() -> None

Write the property.

Source code in repo/hololinked/hololinked/server/http/handlers.py
async def put(self) -> None:
    """Write the property."""
    if await self.is_method_allowed("PUT"):
        self.set_custom_default_headers()
        await self.handle_through_thing(Operations.writeproperty)
    self.finish()

delete async

delete() -> None

Delete the property.

Source code in repo/hololinked/hololinked/server/http/handlers.py
async def delete(self) -> None:
    """Delete the property."""
    if await self.is_method_allowed("DELETE"):
        self.set_custom_default_headers()
        await self.handle_through_thing(Operations.deleteproperty)
    self.finish()

ActionHandler

Bases: RPCHandler

handles action requests.

Source code in repo/hololinked/hololinked/server/http/handlers.py
class ActionHandler(RPCHandler):
    """handles action requests."""

    async def get(self) -> None:
        """Invoke the action, or fetch the reply of an earlier no-block invocation."""
        if await self.is_method_allowed("GET"):
            self.set_custom_default_headers()
            if self.message_id is not None:
                await self.handle_no_block_response()
            else:
                await self.handle_through_thing(Operations.invokeaction)
        self.finish()

    async def post(self) -> None:
        """Invoke the action."""
        if await self.is_method_allowed("POST"):
            self.set_custom_default_headers()
            await self.handle_through_thing(Operations.invokeaction)
        self.finish()

    async def put(self) -> None:
        """Invoke the action."""
        if await self.is_method_allowed("PUT"):
            self.set_custom_default_headers()
            await self.handle_through_thing(Operations.invokeaction)
        self.finish()

    async def delete(self) -> None:
        """Invoke the action."""
        if await self.is_method_allowed("DELETE"):
            self.set_custom_default_headers()
            await self.handle_through_thing(Operations.invokeaction)
        self.finish()

Functions

get async

get() -> None

Invoke the action, or fetch the reply of an earlier no-block invocation.

Source code in repo/hololinked/hololinked/server/http/handlers.py
async def get(self) -> None:
    """Invoke the action, or fetch the reply of an earlier no-block invocation."""
    if await self.is_method_allowed("GET"):
        self.set_custom_default_headers()
        if self.message_id is not None:
            await self.handle_no_block_response()
        else:
            await self.handle_through_thing(Operations.invokeaction)
    self.finish()

post async

post() -> None

Invoke the action.

Source code in repo/hololinked/hololinked/server/http/handlers.py
async def post(self) -> None:
    """Invoke the action."""
    if await self.is_method_allowed("POST"):
        self.set_custom_default_headers()
        await self.handle_through_thing(Operations.invokeaction)
    self.finish()

put async

put() -> None

Invoke the action.

Source code in repo/hololinked/hololinked/server/http/handlers.py
async def put(self) -> None:
    """Invoke the action."""
    if await self.is_method_allowed("PUT"):
        self.set_custom_default_headers()
        await self.handle_through_thing(Operations.invokeaction)
    self.finish()

delete async

delete() -> None

Invoke the action.

Source code in repo/hololinked/hololinked/server/http/handlers.py
async def delete(self) -> None:
    """Invoke the action."""
    if await self.is_method_allowed("DELETE"):
        self.set_custom_default_headers()
        await self.handle_through_thing(Operations.invokeaction)
    self.finish()

RWMultiplePropertiesHandler

Bases: ActionHandler

handles read-write of multiple properties via an action.

Source code in repo/hololinked/hololinked/server/http/handlers.py
class RWMultiplePropertiesHandler(ActionHandler):
    """handles read-write of multiple properties via an action."""

    def initialize(  # ty: ignore[invalid-method-override]
        self,
        resource: ActionAffordance,
        config: RuntimeConfig,
        logger: structlog.stdlib.BoundLogger,
        thing: Thing,
        metadata: HandlerMetadata | None = None,
        **kwargs,
    ) -> None:
        """Set up the handler with the affordances that read and write multiple properties."""
        self.read_properties_resource = kwargs.get("read_properties_resource", None)
        self.write_properties_resource = kwargs.get("write_properties_resource", None)
        return super().initialize(resource, config, logger, thing, metadata)

    async def get(self) -> None:
        """Read multiple properties, or fetch the reply of an earlier no-block read."""
        if await self.is_method_allowed("GET"):
            self.set_custom_default_headers()
            self.resource = self.read_properties_resource
            if self.message_id is not None:
                await self.handle_no_block_response()
            else:
                await self.handle_through_thing(Operations.invokeaction)
        self.finish()

    async def post(self) -> None:
        """Reject the request with 405, multiple properties are written with PUT or PATCH."""
        if await self.is_method_allowed("POST"):
            self.set_status(405, "method not allowed, PUT instead")
        self.finish()

    async def put(self) -> None:
        """Write multiple properties."""
        if await self.is_method_allowed("PUT"):
            self.set_custom_default_headers()
            self.resource = self.write_properties_resource
            await self.handle_through_thing(Operations.invokeaction)
        self.finish()

    async def patch(self) -> None:  # ty: ignore[invalid-method-override]
        """Write multiple properties."""
        if await self.is_method_allowed("PATCH"):
            self.set_custom_default_headers()
            self.resource = self.write_properties_resource
            await self.handle_through_thing(Operations.invokeaction)
        self.finish()

Functions

initialize

initialize(resource: ActionAffordance, config: RuntimeConfig, logger: BoundLogger, thing: Thing, metadata: HandlerMetadata | None = None, **kwargs) -> None

Set up the handler with the affordances that read and write multiple properties.

Source code in repo/hololinked/hololinked/server/http/handlers.py
def initialize(  # ty: ignore[invalid-method-override]
    self,
    resource: ActionAffordance,
    config: RuntimeConfig,
    logger: structlog.stdlib.BoundLogger,
    thing: Thing,
    metadata: HandlerMetadata | None = None,
    **kwargs,
) -> None:
    """Set up the handler with the affordances that read and write multiple properties."""
    self.read_properties_resource = kwargs.get("read_properties_resource", None)
    self.write_properties_resource = kwargs.get("write_properties_resource", None)
    return super().initialize(resource, config, logger, thing, metadata)

get async

get() -> None

Read multiple properties, or fetch the reply of an earlier no-block read.

Source code in repo/hololinked/hololinked/server/http/handlers.py
async def get(self) -> None:
    """Read multiple properties, or fetch the reply of an earlier no-block read."""
    if await self.is_method_allowed("GET"):
        self.set_custom_default_headers()
        self.resource = self.read_properties_resource
        if self.message_id is not None:
            await self.handle_no_block_response()
        else:
            await self.handle_through_thing(Operations.invokeaction)
    self.finish()

post async

post() -> None

Reject the request with 405, multiple properties are written with PUT or PATCH.

Source code in repo/hololinked/hololinked/server/http/handlers.py
async def post(self) -> None:
    """Reject the request with 405, multiple properties are written with PUT or PATCH."""
    if await self.is_method_allowed("POST"):
        self.set_status(405, "method not allowed, PUT instead")
    self.finish()

put async

put() -> None

Write multiple properties.

Source code in repo/hololinked/hololinked/server/http/handlers.py
async def put(self) -> None:
    """Write multiple properties."""
    if await self.is_method_allowed("PUT"):
        self.set_custom_default_headers()
        self.resource = self.write_properties_resource
        await self.handle_through_thing(Operations.invokeaction)
    self.finish()

patch async

patch() -> None

Write multiple properties.

Source code in repo/hololinked/hololinked/server/http/handlers.py
async def patch(self) -> None:  # ty: ignore[invalid-method-override]
    """Write multiple properties."""
    if await self.is_method_allowed("PATCH"):
        self.set_custom_default_headers()
        self.resource = self.write_properties_resource
        await self.handle_through_thing(Operations.invokeaction)
    self.finish()