UniteLabs
How-to

Emit lineage events

Record material and labware movement from a Python workflow.

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.

workflow.py
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:

ArgumentQuestionTransfer example
actorWhich device or service performed it?liquid_handler
operationWhat happened?ASPIRATE
inputsWhich entity did material come from?source_plate:A1
outputsWhich entity received it?channel:0
extrasWhich additional facts describe the event?volume and channel
event_typeWhich 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:

FieldPurpose
roleResource category used by the Platform viewer.
nameReadable label shown in the viewer.
typeLabware type and fallback category when role is absent.
barcodeOptional identity shown on the resource and accepted by trace search.

The current viewer recognizes these role values:

  • wellplate
  • tip_rack
  • carrier
  • reservoir

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:

ResourceExample
Plate wellsource_plate:A1
Tip-rack slottips_300:A1
Carrier sitecarrier:2
Channelchannel:0

Use A1-style addresses for plate wells and tip-rack slots so they map to the Platform grids.

The current viewer treats references other than 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:

CategoryUse for
lineageMaterial, tip, and labware movement.
identificationAn identity associated with an entity, such as a barcode.
measurementA measured value associated with an entity.
auditDevice 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:

ValueSerialized valueMeaning
ASPIRATEaspirateDraw liquid into a channel.
DISPENSEdispenseExpel liquid from a channel.
PICK_UP_TIPpick_up_tipMount a tip on a channel.
PUT_DOWN_TIPput_down_tipReturn a tip to a rack.
DISCARD_TIPdiscard_tipDiscard a tip to waste.
MOVE_LABWAREmove_labwareMove labware between positions.
LOAD_LABWAREload_labwareBring labware into the system.
UNLOAD_LABWAREunload_labwareRemove 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 @phase function exits.
  • The @workflow function 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:

FieldMeaning
totalPreviously unflushed events carrying a flow or task run ID.
deliveredEqual to total after complete success; otherwise 0.
pendingCaptured events retained for another attempt.
attemptsDelivery attempts made by this call.
errorException from the last failed request, if any.
okTrue 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.

Events stay in memory until a flush completes. If the process or pod crashes, you lose every event that wasn't delivered yet.

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.
  • FlushResult does not report individual events rejected inside a successful API response.

Last updated