Emit lineage events
Lineage events record what happened to samples and labware during a workflow run. You describe the labware in your workcell, emit an event after each operation, and the SDK sends both to the UniteLabs Platform.
For the model and Platform views, see Data lineage. For API access, see Read lineage events.
Prerequisites
- UniteLabs SDK ≥ 0.15.0
- A workflow executed on the UniteLabs Platform
A lineage event needs an active run context so the SDK can link it to the run. Code inside @workflow, @phase, and @step functions has that context.
Record a transfer
The example below registers the labware during workflow initialization, then records liquid moving from source_plate:A1, through channel:0, to target_plate:B1 in aspirate and dispense steps within a transfer phase.
from unitelabs.sdk import Audit, Operation, get_context, phase, step, workflow
SOURCE_PLATE_ID = "source_plate"
TARGET_PLATE_ID = "target_plate"
@phase(name="Initialize workflow")
async def initialize_workflow() -> None:
ctx = get_context()
ctx.state["labware_context"] = {
SOURCE_PLATE_ID: {
"name": "Source plate",
"type": "Axygen_96_DW_2mL",
"role": "wellplate",
"barcode": "SOURCE-001",
},
TARGET_PLATE_ID: {
"name": "Target plate",
"type": "Eppendorf_96_PCR",
"role": "wellplate",
"barcode": "TARGET-001",
},
}
@step(name="Aspirate liquid")
async def aspirate_liquid() -> None:
Audit.emit(
actor="liquid_handler",
operation=Operation.ASPIRATE,
inputs={"entity_id": f"{SOURCE_PLATE_ID}:A1"},
outputs={"entity_id": "channel:0"},
extras={"volume_ul": 300, "channel": 0},
)
@step(name="Dispense liquid")
async def dispense_liquid() -> None:
Audit.emit(
actor="liquid_handler",
operation=Operation.DISPENSE,
inputs={"entity_id": "channel:0"},
outputs={"entity_id": f"{TARGET_PLATE_ID}:B1"},
extras={"volume_ul": 300, "channel": 0},
)
@phase(name="Transfer liquid")
async def transfer_liquid() -> None:
await aspirate_liquid()
await dispense_liquid()
@workflow(name="Record transfer lineage")
async def record_transfer_lineage() -> None:
await initialize_workflow()
await transfer_liquid()
This example records lineage only. In a hardware workflow, call Audit.emit immediately after the corresponding operation succeeds.
Run the workflow on the Platform, then open its Lineage tab. Well History shows the events for the 2 wells and the channel. Lineage Graph shows the path from the source well to the target well.
The SDK sends the workspace and events when each phase exits, so this example doesn't need a manual flush.
Understand an event
Each Audit.emit call answers these questions:
| Argument | Question | Transfer example |
|---|---|---|
actor | Which device or service performed it? | liquid_handler |
operation | What happened? | ASPIRATE |
inputs | Which entity did material come from? | source_plate:A1 |
outputs | Which entity received it? | channel:0 |
extras | Which additional facts describe the event? | volume and channel |
event_type | Which event category does it belong to? | lineage |
event_type defaults to lineage. The SDK adds a unique event ID, timestamp, schema version, and active Prefect flow or task run IDs.
Audit.emit appends the event to the run context synchronously. Outside an active context, it returns without recording anything.
Register workspace context
Events use compact entity IDs such as source_plate:A1. Workspace context tells the Platform what source_plate means and gives it a readable label.
Register every referenced container before the first phase or workflow boundary:
ctx = get_context()
ctx.state["labware_context"] = {
"source_plate": {
"name": "Source plate",
"type": "Axygen_96_DW_2mL",
"role": "wellplate",
"barcode": "SOURCE-001",
}
}
A descriptor supports these fields:
| Field | Purpose |
|---|---|
role | Resource category used by the Platform viewer. |
name | Readable label shown in the viewer. |
type | Labware type and fallback category when role is absent. |
barcode | Optional identity shown on the resource and accepted by trace search. |
The current viewer recognizes these role values:
wellplatetip_rackcarrierreservoir
There's no tube role. The API stores a tube descriptor, but the viewer shows it as generic labware.
Refer to child resources
Labware children use the <parent>:<key> identifier format:
| Resource | Example |
|---|---|
| Plate well | source_plate:A1 |
| Tip-rack slot | tips_300:A1 |
| Carrier site | carrier:2 |
| Channel | channel:0 |
Use A1-style addresses for plate wells and tip-rack slots so they map to the Platform grids.
channel:<n> as labware. It has no resource type for an out-of-system destination.Workspace lifecycle
The Platform receives the workspace once per run, at the first phase or workflow boundary. Changes you make to labware_context after that never reach it: later calls to Audit.flush_workspace() do nothing. Declare every container the run will touch, including ones loaded mid-run, before the first phase starts. To attach a barcode you scan during the run, emit an identification event instead (see Choose an event category).
Choose an event category
The SDK defines 4 standard categories:
| Category | Use for |
|---|---|
lineage | Material, tip, and labware movement. |
identification | An identity associated with an entity, such as a barcode. |
measurement | A measured value associated with an entity. |
audit | Device state, operator decisions, or run milestones. |
The API also accepts custom category strings. Hamilton load_carrier and unload_carrier steps, for example, use a custom error category when the operation fails. Consumers can filter events by type.
To record an identification, add this call inside the example workflow after a successful scan:
Audit.emit(
actor="barcode_reader",
operation=Operation.LOAD_LABWARE,
inputs={},
outputs={"entity_id": SOURCE_PLATE_ID, "barcode": "SOURCE-001"},
event_type="identification",
)
Choose an operation
Audit.emit expects one of these Operation enum members:
| Value | Serialized value | Meaning |
|---|---|---|
ASPIRATE | aspirate | Draw liquid into a channel. |
DISPENSE | dispense | Expel liquid from a channel. |
PICK_UP_TIP | pick_up_tip | Mount a tip on a channel. |
PUT_DOWN_TIP | put_down_tip | Return a tip to a rack. |
DISCARD_TIP | discard_tip | Discard a tip to waste. |
MOVE_LABWARE | move_labware | Move labware between positions. |
LOAD_LABWARE | load_labware | Bring labware into the system. |
UNLOAD_LABWARE | unload_labware | Remove labware from the system. |
The API stores any string in its operation field, but Audit.emit serializes an Operation enum member. Choose the closest enum value and add domain-specific detail to extras when the enum has no exact match.
Most Liquid Handling SDK methods don't emit lineage automatically. The current exceptions are Hamilton load_carrier and unload_carrier steps, which emit LOAD_LABWARE and UNLOAD_LABWARE events on success and failure. Failed operations use event_type="error".
Automatic delivery
The SDK saves the run context and tries to send its workspace and events when:
- A
@phasefunction exits. - The
@workflowfunction exits.
Both boundaries run after success or failure. A @step exit does not trigger a save or flush.
Boundary delivery is non-strict, with a 30-second deadline. If delivery fails, the SDK logs the failure and the workflow keeps its own result or exception. Pending events stay in memory, so a later flush in the same process can retry them.
Flush manually
Call Audit.flush() when events must appear during a long-running phase or when you need to inspect the request result at a specific point:
from unitelabs.sdk import Audit, get_logger
logger = get_logger(__name__)
result = await Audit.flush()
if not result.ok:
logger.warning(
"Lineage delivery incomplete: %d events pending after %d attempts: %s",
result.pending,
result.attempts,
result.error,
)
A successful manual flush advances the delivery position, so the next flush only sends events added after it.
Advanced delivery controls
Delivery result
Audit.flush() returns a FlushResult:
| Field | Meaning |
|---|---|
total | Previously unflushed events carrying a flow or task run ID. |
delivered | Equal to total after complete success; otherwise 0. |
pending | Captured events retained for another attempt. |
attempts | Delivery attempts made by this call. |
error | Exception from the last failed request, if any. |
ok | True when request failures left no captured events pending. |
FlushResult reports committed request completion. A failed flush may have sent some batches, but the SDK leaves all captured events pending so the next flush can resend them for server-side deduplication. The API can also return a successful response and reject an individual event when it cannot resolve that event to a run. The current SDK does not include the response's accepted, duplicates, or rejected values in FlushResult.
Fail on a request error
Use strict mode when a failed upload request should fail the run:
from unitelabs.sdk import Audit
await Audit.flush(strict=True, max_attempts=6, backoff=5)
Strict mode raises LineageDeliveryError after the retry budget is exhausted. The exception carries the FlushResult in its result attribute. Strict mode does not raise for an individual event rejected inside a successful API response.
Timeouts and retries
Each client request uses a default (connect, read) timeout of (10.0, 60.0) seconds:
from unitelabs.sdk import AsyncApiClient, Audit
async with AsyncApiClient(timeout=(10.0, 60.0)) as client:
result = await Audit.flush(client=client, max_attempts=3, backoff=1.0)
Pass timeout=None to wait without a client timeout. max_attempts includes the first request. Retries use linear backoff.
Keep in mind
- The workspace is sent once per run context.
- The Platform viewer has no first-class tube or out-of-system resource type.
- Events without a Prefect flow or task run ID are skipped during delivery.
- Automatic boundary delivery logs and suppresses request failures.
FlushResultdoes not report individual events rejected inside a successful API response.
Last updated