Skip to content

vllm_omni.entrypoints.duplex.session_attachment

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.

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]],
) -> 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.

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

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,
) -> 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.

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