Skip to content

vllm_omni.entrypoints.duplex.session_attachment

DuplexDetachedAttachment dataclass

The transport a detach dropped, with the engine lease generation it was serving.

attachment_generation instance-attribute

attachment_generation: int

lease_generation instance-attribute

lease_generation: int

session_id instance-attribute

session_id: str

DuplexEventJournal

last_sequence property

last_sequence: int

overflowed property

overflowed: bool

retained_bytes property

retained_bytes: int

acknowledge

acknowledge(sequence: int) -> int

prune

prune(now: float | None = None) -> int

record

record(payload: Mapping[str, object]) -> JournalEntry

replay_after

replay_after(sequence: int) -> tuple[JournalEntry, ...]

DuplexJournalGapError

Bases: RuntimeError

DuplexJournalOverflowError

Bases: RuntimeError

DuplexResumeCredential dataclass

token_digest instance-attribute

token_digest: bytes

from_token classmethod

from_token(token: ResumeToken) -> DuplexResumeCredential

rotate

rotate() -> ResumeToken

verify

verify(plaintext: str) -> bool

DuplexSessionAttachmentCreated dataclass

attachment_generation instance-attribute

attachment_generation: int

resume_token class-attribute instance-attribute

resume_token: ResumeToken = field(repr=False)

session_id instance-attribute

session_id: str

DuplexSessionAttachmentRegistry

acknowledge async

acknowledge(session_id: str, sequence: int) -> int

authenticate_resume async

authenticate_resume(
    session_id: str,
    *,
    resume_token: str,
    last_received_server_event_seq: int,
) -> None

Validate transport credentials before any engine resume control.

begin_resume async

begin_resume(session_id: str) -> None

Claim that a resume of session_id is about to land and activate.

Held from before the engine lease resume until end_resume. While any claim is open, a lease generation that loses its owner is parked for the next activation (see settle_lease_generation) instead of being handed to the currently attached connection, which that activation is about to replace.

close async

close(session_id: str) -> DuplexTransportAttachment | None

create async

create(
    session_id: str,
    *,
    send: Callable[[dict[str, object]], Awaitable[None]],
    close: Callable[[str], Awaitable[None]],
    lease_generation: int = 0,
) -> DuplexSessionAttachmentCreated

detach async

detach(
    session_id: str,
    *,
    attachment_generation: int | None = None,
) -> bool

Drop the current transport; the engine lease owns the disconnect grace.

attachment_generation names the connection asking to detach, so a socket that already lost a takeover cannot detach the winner. None means "whichever connection is attached right now" and is for callers that only know the session (the outbound pump, whose send just failed).

Returns whether this call is the one that detached: an already-detached session answers False so a second disconnect signal for the same socket cannot restart the engine's disconnect grace window.

end_resume async

end_resume(session_id: str) -> int | None

Release a begin_resume claim, whether or not the resume activated.

Returns a lease generation the caller must detach itself: one that was parked for an activation that never came, with nobody attached to take it over. None when nothing is owed.

has_attachment async

has_attachment(session_id: str) -> bool

Whether some connection is attached right now, whoever it is.

A resume that failed to activate has no generation of its own, so it cannot ask is_current_attachment. What it needs to know before rolling the engine lease back into its disconnect grace is only whether it would be rolling back somebody else's live attachment.

is_current_attachment async

is_current_attachment(
    session_id: str, attachment_generation: int
) -> bool

release_attachment async

release_attachment(
    session_id: str,
    *,
    attachment_generation: int | None = None,
) -> DuplexDetachedAttachment | None

detach that also hands back the lease generation the dropped connection was serving.

Captured under the same lock as the attachment: the shared session handle may already carry a newer generation from a resume that has not activated yet, and the engine detach that follows must be fenced on the dropped connection's own lease, not on that one.

resume async

resume(
    session_id: str,
    *,
    resume_token: str,
    last_received_server_event_seq: int,
    send: Callable[[dict[str, object]], Awaitable[None]],
    close: Callable[[str], Awaitable[None]],
    activation_payload_factory: Callable[
        [ResumeToken, int], Mapping[str, object]
    ]
    | None = None,
    lease_generation: int | None = None,
) -> DuplexSessionResumeResult

send_event async

send_event(
    session_id: str,
    payload: Mapping[str, object],
    *,
    journal: bool = True,
    on_accepted: Callable[[], None] | None = None,
) -> JournalEntry | None

Sequence and dispatch one event to the current attachment.

The per-session lock keeps wire order equal to journal order without serializing unrelated sessions. A detached session still records replayable events, but has no transport side effect.

on_accepted runs synchronously once the event is journaled, or after a successful transport send when journaling is disabled. A later send failure cannot undo acceptance into the journal. Detached, non-journaled events do not invoke it. The callback must not raise.

settle_lease_generation async

settle_lease_generation(
    session_id: str, lease_generation: int
) -> int | None

Find an owner for the engine lease generation of a resume its connection never served.

The engine already bumped its lease to lease_generation. The owner is, in order: the resume waiting to activate (it will serve the session and absorbs the generation when it activates); else the connection attached right now (it keeps serving, and its own disconnect must be able to detach this lease); else nobody, in which case the generation is returned and the caller puts it back into disconnect grace. Generations only move forward, so an older orphan never displaces a newer one an owner already holds.

DuplexSessionResumeResult dataclass

attachment_generation instance-attribute

attachment_generation: int

replaced_attachment class-attribute instance-attribute

replaced_attachment: DuplexTransportAttachment | None = None

replay_entries class-attribute instance-attribute

replay_entries: tuple[JournalEntry, ...] = ()

resume_token class-attribute instance-attribute

resume_token: ResumeToken = field(repr=False)

session_id instance-attribute

session_id: str

DuplexTransportAttachment dataclass

close class-attribute instance-attribute

close: Callable[[str], Awaitable[None]] = field(repr=False)

generation instance-attribute

generation: int

send class-attribute instance-attribute

send: Callable[[dict[str, object]], Awaitable[None]] = (
    field(repr=False)
)

InvalidResumeTokenError

Bases: RuntimeError

JournalEntry dataclass

created_monotonic instance-attribute

created_monotonic: float

encoded_bytes instance-attribute

encoded_bytes: int

payload class-attribute instance-attribute

payload: Mapping[str, object] = field(repr=False)

sequence instance-attribute

sequence: int

ResumeToken dataclass

plaintext class-attribute instance-attribute

plaintext: str = field(repr=False)

generate classmethod

generate() -> ResumeToken