A FastAPI app that communicates with a generic Temporal and streams updates via websocket.
Project description
Fast Temporal
A Python package that provides a FastAPI application with WebSocket support for real-time communication with Temporal workflows. This package enables streaming updates from Temporal workflows to clients through WebSocket connections.
Features
- FastAPI server with WebSocket support
- Real-time workflow status updates
- Generic Temporal workflow base class with activity scheduling functions and query handler.
- Support for multiple workflows.
- Environment-based configuration
- CORS support
- Structured logging
Installation
pip install fast-temporal
Configuration
Create a .env file in your project root with the following variables:
TEMPORAL_CLIENT=localhost:7233
POLLING_INTERVAL=0.5
ALLOWED_ORIGINS=*
FASTAPI_HOST=0.0.0.0
FASTAPI_PORT=8000
FASTAPI_RELOAD=true
Usage
GenericTemporalWorkflow Class
The GenericTemporalWorkflow class provides a robust foundation for building Temporal workflows with built-in activity scheduling, state management, and real-time status updates. This can be extended by every workflow involved in the implementation.
Key Functions
-
Activity Scheduling
async def schedule_activity( self, activity_name: str, callback: Optional[Callable] = None, args: List[Any] = None, kwargs: Dict[str, Any] = None, timeout: int = 60, retry_policy: Optional[RetryPolicy] = None ) -> Any
- Schedules and executes a Temporal activity
- Supports optional callback for result processing
- Configurable timeout and retry policy
- Returns activity result or None if failed
-
Workflow Result Management
def set_workflow_result(self, result: Any, status: str = "Done") -> None
- Sets the final workflow result
- Updates workflow status (default: "Done")
- Triggers workflow completion
-
Query Handlers
@workflow.query async def get_current_activity(self) -> Dict[str, str]
- Returns current activity name, status, and ID
- Used for real-time status tracking
- Format:
{"current_activity": name, "status": status, "activity_id": id}
@workflow.query async def get_activity_result(self, activity_id: str) -> Any
- Retrieves result of a completed activity
- Used for accessing activity outputs
- Returns None if activity not found
@workflow.query async def get_callback_result(self, activity_id: str) -> Any
- Retrieves the result of a callback function associated with a completed activity
- Used for accessing the output of a callback after an activity finishes
- Arguments:
activity_id(str): The unique identifier of the activity whose callback result you want to fetch
- Returns:
- The result of the callback, or
Noneif not found
- The result of the callback, or
-
State Management
def set_state(self, key: str, value: Any) -> None def get_state(self, key: str, default: Any = None) -> Any
- Stores and retrieves workflow state
- Useful for sharing data between activities
- Supports custom key-value pairs
Usage Example
@workflow.defn
class MyWorkflow(GenericTemporalWorkflow):
@workflow.run
async def run(self, input_data: Dict[str, Any]):
# Schedule an activity with callback
result = await self.schedule_activity(
"process_data",
args=[input_data],
callback=self.handle_result
)
return result
async def handle_result(self, result):
# Process activity result
processed = do_something(result)
# Set workflow result
self.set_workflow_result(processed)
return processed
Starting the Server
After installation, you can start the FastAPI server using the provided script:
fast-temporal-run
Optional command-line arguments:
--host: Server host (default: from .env)--port: Server port (default: from .env)--reload: Enable auto-reload (default: from .env)
Example:
fast-temporal-run --host 127.0.0.1 --port 8080 --reload
WebSocket Communication
Connect to the WebSocket endpoint at /ws/{user_id} where user_id is a unique identifier for your client.
Example client connection:
ws_url = f"ws://localhost:8000/ws/{user_id}"
async with websockets.connect(ws_url) as ws:
data={"args": {"prompt": prompt, "user_id": user_id}, "origin": "streamlit_ui"}
await ws.send(json.dumps(data))
while True:
response = await ws.recv()
data = json.loads(response)
WebSocket Communication
Connect to the WebSocket endpoint at /ws/{user_id} where user_id is a unique identifier for your client.
Sending Data to Start a Workflow
To start a workflow, you must send a JSON message over the WebSocket with a specific structure. The backend expects a dictionary containing args, origin, and workflow keys.
It is crucial that the keys in the dictionary are named exactly as shown below, as the backend uses these specific keys to process the request.
{
"args": {
"prompt": "your prompt here",
"user_id": "unique_user_id"
},
"origin": "streamlit_ui",
"workflow": {
"workflow_name": "TestWorkflow",
"workflow_task_queue": "test-task-queue",
"start_signal_function": "handle_llm_request"
}
}
Key-Value Explanations:
args(dict): A dictionary containing the arguments to be passed to your workflow's start signal function. The structure of this dictionary is dependent on your specific workflow's requirements.origin(str): A string that identifies the client application sending the request (e.g.,"streamlit_ui"). This helps in logging and routing. For messages received from temporal, this value will be given out as"temporal"workflow(dict): A dictionary that provides the necessary metadata to identify and start the correct Temporal workflow.workflow_name(str): The name of the workflow class that you want to execute. This must match a workflow registered with your Temporal worker.workflow_task_queue(str): The task queue that the workflow will be scheduled on. Ensure your Temporal worker is listening to this queue.start_signal_function(str): The name of the signal method within your workflow that will be triggered to start the execution. Theargsdictionary will be passed as an argument to this method.
Workflow Updates
The server will send real-time updates about workflow activities in the following format:
{
"origin": "temporal",
"message": "activity_name: status",
"status": "Running|Completed|Failed|Done"
}
Final Response
The final response will be sent when the status becomes Done. Done indicates that the workflow is complete and will be set when the workflow result is set.
Therefore, the workflow result must be set using the set_workflow_result method, ideally in your CALLBACK functions of your final activity. This action sets the workflow result, marks the current status as Done, and completes the workflow.
Optionally, you can set the status as any other status using the set_workflow_result(result,"Failed").
The workflow will still be completed if you dont want retries. The example contains both the scenarios - Activity task failures that will show up in the Temporal workflow UI and Activity task failure that are handled and won't show up in the UI.
The result of the final activity, that was run will be sent through the websocket in the following JSON format:
{
"origin": "temporal",
"message": final_activity_result
"status": "Done"
}
NOTE:
All the activity results, if needed, can be retrieved using the get_activity_result query handler. To fetch the result of callback, please use the get_callback_result query handler, the argument supplied should be the activity ID for which the callback result is needed.
Example Application
The package includes an example application in the examples/example_app directory that demonstrates the integration of FastAPI, Temporal, and Streamlit. The example shows how to:
- Define Temporal activities and workflows
- Set up a Temporal worker
- Create a FastAPI server with WebSocket support
- Build a Streamlit UI that communicates with the workflow
In this example, we are taking in input from the user, where he uploads a txt file. Then according to user instructions, we generate the content by calling an LLM(Activity 1), write it into the text file (Activity 2), and create an audio file(Activity 3).
Running the Example
- Start the Temporal worker:
python examples/example_app/temporal_worker.py
- Start the FastAPI server:
fast-temporal-run
- Launch the Streamlit UI:
streamlit run examples/example_app/streamlit_ui.py
The example demonstrates:
- Real-time workflow status updates
- Activity scheduling and management
- WebSocket communication
- Streamlit UI integration
Package Structure
fast_temporal/
├── fast_temporal/
│ ├── config/
│ ├── workflow/
│ └── api/
│
├── examples/
│ └── example_app/
│ └── activities/
│
├── pyproject.toml
└── README.md
Key Components
-
Config Module
- Environment variable management
- Logger configuration
- Configuration validation
-
Workflow Module
- Generic Temporal workflow base class
- Activity scheduling and management
- State management
- Error handling
-
API Module
- FastAPI application setup
- WebSocket connection management
- Temporal client integration
- Real-time status updates
Dependencies
- fastapi
- uvicorn[standard]
- python-dotenv
- temporalio
- websockets
Contributing
- Fork the repository
- Create your feature branch
- Commit your changes
- Push to the branch
- Create a new Pull Request
Project details
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file fast_temporal-0.2.1.1.tar.gz.
File metadata
- Download URL: fast_temporal-0.2.1.1.tar.gz
- Upload date:
- Size: 14.0 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.1.0 CPython/3.11.9
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
31cac0fae3176254dbb633babd899fa6b7331c8f3bcd88d3204f6bcd4026c500
|
|
| MD5 |
c6e387ece1a6c969a734e47cec52c3b2
|
|
| BLAKE2b-256 |
ca10474ce0467e1a3c283bb956565c337f2abbbbaad82aee5e7ea2404bb506b0
|
File details
Details for the file fast_temporal-0.2.1.1-py3-none-any.whl.
File metadata
- Download URL: fast_temporal-0.2.1.1-py3-none-any.whl
- Upload date:
- Size: 12.0 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.1.0 CPython/3.11.9
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
444fe2c2d76f01cadffb77edbfa3f1db26a53403d9cbed2c66c0338a33f604d1
|
|
| MD5 |
2ab1dff0fbb4f6d7d65a5694b004fc7a
|
|
| BLAKE2b-256 |
12b52f9dc8e9cd5760d5e223264fffdbcedadfe108a482ce949e1a019bd132bf
|