Skip to main content

is-wire-sea

Middleware AMQP em Python para a arquitetura IS. A distribuição no PyPI se chama is-wire-sea; o pacote Python estável continua sendo is_wire.

Instalação

python -m pip install is-wire-sea

Há suporte para Python 3.10 a 3.14. Instale um exportador OTLP com:

python -m pip install 'is-wire-sea[tracing]'

Broker

Implantações de produção devem fixar o RabbitMQ pela versão de patch e pelo digest da imagem. A base de desenvolvimento é o RabbitMQ 4.3.1:

docker run --rm -p 5672:5672 \
  rabbitmq:4.3.1@sha256:6a46d2aef889d2a8cc28ac91b4a1ca0116a4a151de10cada219ee7685dc01c5b

RabbitMQ 3.7.6 é apenas um alvo de migração. Mova as implantações existentes da versão 3.7.6 para um cluster 4.3 novo usando uma migração blue-green; não reutilize diretamente o diretório de dados antigo.

Pub/sub

A API existente de entrega no máximo uma vez é preservada:

from is_wire.core import Channel, Message, Subscription

with Channel("amqp://guest:guest@localhost:5672") as channel:
    subscription = Subscription(channel)
    subscription.subscribe("Camera.Config")

    channel.publish(Message(content=b"hello"), topic="Camera.Config")
    received = channel.consume(timeout=1.0)

Assinaturas anônimas são exclusivas e temporárias. Assinaturas nomeadas são duráveis, compartilhadas entre réplicas e expiram após cinco minutos sem uso.

Fluxos de imagens em tempo real

Use StreamSubscription quando quadros antigos puderem ser descartados. Réplicas com o mesmo grupo compartilham uma fila, então cada quadro é processado por uma réplica em vez de ser copiado para todas elas.

from is_msgs.image_pb2 import Image
from is_wire.core import Channel, Message, StreamSubscription


def process(message):
    image = message.unpack(Image)
    # image.data já contém bytes JPEG, WebP ou PNG.


with Channel("amqp://guest:guest@localhost:5672") as channel:
    stream = StreamSubscription(channel, group="detector")
    stream.subscribe("Camera.*.Frame")
    stream.run(process)

Publique um quadro com semântica de fluxo:

image = Image(data=encoded_jpeg_bytes)
channel.publish_stream(Message(content=image), topic="Camera.0.Frame")

O preset de fluxo usa confirmações manuais, prefetch 1, fila com tamanho um, descarte do item mais antigo quando há overflow, TTL de dois segundos e mensagens não persistentes. Uma exceção no callback rejeita o quadro sem recolocá-lo na fila. Use outro grupo somente quando outra etapa de processamento realmente precisar de sua própria cópia.

Image.data deve conter bytes binários codificados da imagem. Não envie RGB/BGR bruto, a menos que o orçamento de rede permita explicitamente; não codifique o corpo em base64 e não aplique compressão gzip/zstd genérica a um payload JPEG, WebP ou PNG já comprimido.

Use um Channel separado para o tráfego de imagens e para o tráfego de RPC/controle, evitando que quadros grandes atrasem mensagens pequenas de controle.

Protobuf

from google.protobuf.struct_pb2 import Struct
from is_wire.core import ContentType, Message

value = Struct()
value["camera"] = "front"

binary = Message(content=value)
assert binary.content_type == ContentType.PROTOBUF
assert binary.unpack(Struct) == value

json_message = Message(content_type=ContentType.JSON)
json_message.pack(value)

Há suporte para Protobuf 5 a 7. Os corpos binários existentes e as convenções de propriedades AMQP continuam compatíveis com is-wire 1.2.1.

RPC

from google.protobuf.struct_pb2 import Struct
from is_wire.core import Channel
from is_wire.rpc import ServiceProvider


def echo(request, context):
    return request


channel = Channel("amqp://guest:guest@localhost:5672")
provider = ServiceProvider(channel)
provider.delegate("Echo", echo, Struct, Struct)
provider.run()

As filas de serviços RPC usam confirmações manuais e prefetch 16. Uma requisição só é confirmada depois que sua resposta é publicada; por isso, os handlers devem ser idempotentes.

OpenTelemetry

Tracer, TracingInterceptor, Message.inject_tracing e Message.extract_tracing usam OpenTelemetry e propagação B3 com múltiplos cabeçalhos. IDs de trace legados de 64 bits e IDs de 128 bits são aceitos.

from is_wire.core import Message, Tracer

tracer = Tracer()
with tracer.span("publish") as span:
    message = Message(content=b"payload")
    message.inject_tracing(span)

Zipkin

Instale o exportador nativo OpenTelemetry Zipkin JSON v2:

python -m pip install 'is-wire-sea[zipkin]'

Crie um provider compartilhado por processo. Ele agrupa as exportações em segundo plano, aplica amostragem baseada no contexto pai e identifica cada instância de câmera separadamente no Zipkin:

from is_wire.core import Message, ZipkinTracing
from is_wire.rpc import ServiceProvider, TracingInterceptor

tracing = ZipkinTracing(
    service_name="camera-gateway",
    endpoint="http://zipkin:9411/api/v2/spans",
    service_instance_id="camera-5",
    sample_ratio=0.1,
    resource_attributes={"camera.driver": "hikvision"},
)
frame_tracer = tracing.tracer()

provider = ServiceProvider(rpc_channel)
provider.add_interceptor(TracingInterceptor(tracing=tracing))

with frame_tracer.span("camera.frame") as span:
    span.set_attribute("camera.id", "5")
    message = Message(content=image)
    message.inject_tracing(span)
    stream_channel.publish_stream(message, topic="CameraGateway.5.Frame")

# Libere os spans enfileirados durante o encerramento normal do processo.
tracing.shutdown()

Não associe corpos de imagem, credenciais, URLs de câmeras ou identificadores sem limite aos spans. Use uma proporção de amostragem menor que um para vídeo contínuo e defina 1.0 apenas em sessões curtas de diagnóstico. Consumidores continuam o trace do produtor com message.extract_tracing().

Métricas

As métricas do Prometheus são registradas no registry padrão. MetricsInterceptor.start_server() pode expô-las por HTTP. O perfil 2.0 exporta:

  • is_wire_messages_published_total e is_wire_published_bytes_total;
  • is_wire_messages_received_total e is_wire_received_bytes_total;
  • is_wire_stream_processing_seconds, is_wire_stream_frame_age_seconds e is_wire_stream_callback_errors_total;
  • is_wire_reconnections_total;
  • is_wire_rpc_duration_seconds e is_wire_rpc_requests_total, particionados por status.

O conteúdo das imagens e os IDs de correlação nunca são usados como labels de métricas.

TLS e ciclo de vida da conexão

amqps:// habilita TLS com verificação de certificado e hostname. Opções SSL adicionais do py-amqp podem ser passadas em ssl_options. Consumidores se reconectam com backoff exponencial; publicações nunca são repetidas automaticamente, pois a entrega pode já ter ocorrido.

channel = Channel(
    "amqps://user:password@rabbitmq.example:5671/vhost",
    heartbeat=30,
    ssl_options={"ca_certs": "/etc/ssl/certs/cluster-ca.pem"},
)

Desenvolvimento

python -m pip install -e '.[dev,tracing]'
pytest
ruff check src tests
python -m build
twine check --strict dist/*

Regenere o schema wire interno de forma reproduzível com o compilador Protobuf 5.29 fixado:

python -m pip install -e '.[codegen]'
python scripts/generate_wire.py

Execute o perfil de publicador a 30 FPS / consumidor a 10 FPS contra um broker de teste com:

python scripts/validate_stream_profile.py --payload-size 1048576 --frames 300

Execute o mesmo comando em implantações Kubernetes separadas que compartilhem um group para comparar as métricas de entrada/saída do RabbitMQ conforme réplicas são adicionadas.

Consulte MIGRATION.md, COMPATIBILITY.md e CHANGELOG.md antes de atualizar uma implantação existente.

Agradecimentos

A modernização, a revisão de compatibilidade, os testes e o empacotamento no PyPI da série 2.0 foram realizados com assistência do OpenAI Codex.

Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

is_wire_sea-2.0.0.tar.gz (29.0 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

is_wire_sea-2.0.0-py3-none-any.whl (30.6 kB view details)

Uploaded Python 3

File details

Details for the file is_wire_sea-2.0.0.tar.gz.

File metadata

  • Download URL: is_wire_sea-2.0.0.tar.gz
  • Upload date:
  • Size: 29.0 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: uv/0.12.2 {"installer":{"name":"uv","version":"0.12.2","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"20.04","id":"focal","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":null}

File hashes

Hashes for is_wire_sea-2.0.0.tar.gz
Algorithm Hash digest
SHA256 37d9ad6798fdf6d82c33b7f43cd34bcdcc9c46359f55a45fdc6586167ca78bcf
MD5 46eb558896b3ee20a7474d8934ca5316
BLAKE2b-256 236427ebbeb41530da5f989d5e3739d2728279312c4086187abb8e95c2f55559

See more details on using hashes here.

File details

Details for the file is_wire_sea-2.0.0-py3-none-any.whl.

File metadata

  • Download URL: is_wire_sea-2.0.0-py3-none-any.whl
  • Upload date:
  • Size: 30.6 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? No
  • Uploaded via: uv/0.12.2 {"installer":{"name":"uv","version":"0.12.2","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"20.04","id":"focal","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":null}

File hashes

Hashes for is_wire_sea-2.0.0-py3-none-any.whl
Algorithm Hash digest
SHA256 cd3bec67bd45a14a618ae808fa4538fcb920d3bf15f6270b5989cba147407762
MD5 54b0cd070acd61b995cb894ef579e65d
BLAKE2b-256 d938440b7038168a1a748856a6e9c13bd4700b899e71d1752f762e8731017ec9

See more details on using hashes here.

Release history Release notifications | RSS feed

2.0.2

2 files

2.0.1

2 files

This release

2.0.0 This release

2 files

1.3.1

2 files

1.3.0

2 files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page