Fluxos críticos¶
Diagramas de sequência para os fluxos que mais erram na primeira implementação — incluindo a mensageria em tempo real (SSE + Web Push) — junto com as máquinas de estado de Order e Invitation. Cada fluxo aponta os primitivos do SDK envolvidos.
1. Signup público + login¶
sequenceDiagram
autonumber
actor C as Cliente
participant R as auth.router
participant S as UserService
participant U as UserUtils (PasswordUtils)
participant J as JWTUtils
participant DB as Postgres
C->>R: POST /auth/signup {email, password, name}
R->>S: signup(payload)
S->>U: hash(password)
U-->>S: bcrypt hash
S->>DB: INSERT users (email, hash, ...)
DB-->>S: user row
S->>J: encode({sub: user.id}, ttl=ACCESS_TTL)
S->>J: encode({sub: user.id, refresh: true}, ttl=REFRESH_TTL)
S-->>R: {user, access, refresh}
R-->>C: 201 Created
Pontos do SDK:
- Endpoint público —
auth.routernão usaDepends(get_current_user). PasswordUtils.hash(bcrypt) +JWTUtils.encode(HS256).- Falha de email duplicado MUST virar
ConflictException→ handler do SDK responde409com envelope padrão.
2. Convite de membro¶
sequenceDiagram
autonumber
actor A as Admin (OWNER/ADMIN)
actor I as Convidado
participant R as invitations.router
participant S as InvitationService
participant T as generate_opaque_token
participant E as EmailUtils
participant Q as TaskIQ (async email)
participant DB as Postgres
A->>R: POST /organizations/{id}/invitations {email, role}
R->>S: invite(org_id, payload, current_user)
S->>S: assert role ≠ OWNER
S->>S: assert org_member_count < 10
S->>T: generate_opaque_token(48)
T-->>S: (plain, hash)
S->>DB: INSERT invitations (token_hash, expires_at=now+7d, PENDING)
S->>Q: enqueue send_invitation_email(invite.id, plain)
Q-->>E: render_template("invitation.html", {...})
E-->>I: email com link ?token={plain}
S-->>R: invitation
R-->>A: 201
Note over I: 1 dia depois
I->>R: POST /invitations/{plain}/accept (com JWT do convidado)
R->>S: accept(plain, current_user)
S->>S: hash_opaque_token(plain) -> lookup
S->>S: assert convite.email == current_user.email
S->>S: assert not expired & status=PENDING
S->>S: assert org_member_count < 10
S->>DB: BEGIN
S->>DB: INSERT memberships (role=convite.role)
S->>DB: UPDATE invitations SET status=ACCEPTED
S->>DB: COMMIT
S-->>R: membership
R-->>I: 200
Pontos do SDK:
generate_opaque_token(48)retorna par(plain, hash).EmailUtils.render_template("invitation.html", ctx)(v0.24+).- O envio é assíncrono (TaskIQ) — endpoint retorna
201sem esperar SMTP. - Toda a aceitação é uma única transação — membership + status do convite são atomic.
O banco guarda só o hash
generate_opaque_token(48) devolve (plain, hash): o valor plain só existe no email enviado ao convidado, e o banco persiste apenas hash. Na aceitação, o service faz hash_opaque_token(plain) e busca pelo hash — um vazamento da tabela invitations não expõe tokens utilizáveis.
3. Criar produto com variante + imagens¶
sequenceDiagram
autonumber
actor M as Membro (ADMIN+)
participant R as products.router
participant CT as ProductController
participant PS as ProductService
participant VS as VariantService
participant ST as AsyncMinIOClient
participant DB as Postgres
M->>R: POST /products {title, description, variants:[{sku, attrs, price_cents}]}
R->>CT: create_product(payload, org_id, user_id)
CT->>PS: create(org_id, payload)
PS->>DB: BEGIN
PS->>DB: INSERT products (...)
loop pra cada variant
PS->>VS: create_variant(product_id, variant_payload)
VS->>DB: INSERT product_variants (...)
VS->>DB: INSERT price_history (valid_from=now())
end
PS->>DB: COMMIT
PS-->>R: product
Note over M,ST: Upload de imagem (separado)
M->>R: POST /products/{id}/images/presign
R->>ST: presigned_put_url("products/{id}/{uuid}.jpg", 15min)
ST-->>R: {key, url}
R-->>M: {key, url}
M->>ST: PUT bytes direto no MinIO via URL presigned
M->>R: PATCH /products/{id} {image_keys: [...keys]}
R->>PS: attach_images(product_id, keys)
PS->>DB: UPDATE products SET image_keys = ...
PS-->>R: product
R-->>M: 200
Pontos do SDK:
- Criação de produto é transação única — produto + variantes + primeira linha de
PriceHistory. - Imagens não trafegam pela API — cliente faz
PUTdireto no MinIO via URL presigned gerada porAsyncMinIOClient.presigned_put_url(oMinIOUploadStorage.presigned_urlé GET/leitura, não serve pra upload). - Catálogo público lê
image_keyse gera URLs presigned de leitura (TTL 1h).
4. Checkout idempotente¶
sequenceDiagram
autonumber
actor B as Comprador
participant MW as IdempotencyMiddleware
participant R as orders.router
participant OC as OrderController
participant OS as OrderService
participant SS as StockService
participant SSE as orders/{id}/events stream
participant DB as Postgres
participant Q as TaskIQ
B->>MW: POST /orders {cart_id, address}<br/>Idempotency-Key: chk_uuid
MW->>MW: cache lookup (method+path+key)
alt cache hit
MW-->>B: response cacheada (200/201)
else cache miss
MW->>R: forward
R->>OC: checkout(cart_id, address, user)
OC->>OS: create_order(cart, address, user)
OS->>DB: BEGIN
OS->>DB: SELECT cart FOR UPDATE
OS->>OS: assert cart.user == user & status=OPEN
OS->>SS: reserve(items)
loop pra cada item
SS->>DB: assert balance(variant) >= qty
SS->>DB: INSERT stock_movements (kind=RESERVATION)
end
OS->>DB: INSERT orders (status=PENDING, idem_key)
OS->>DB: INSERT order_items (...)
OS->>DB: UPDATE carts SET status=CONVERTED
OS->>DB: COMMIT
OS->>Q: enqueue notify_seller(order.id)
OS->>SSE: publish {order_id, status: PENDING}
OS-->>R: order
R-->>MW: 201 (body completo)
MW->>MW: store response under key
MW-->>B: 201 Created
end
Pontos do SDK:
IdempotencyMiddlewarecobre o endpoint sem o handler precisar saber.- Reserva de estoque é dentro da mesma transação do
INSERTdo pedido. Falha em qualquer item aborta tudo. - A
SSEnotifica o stream (cliente do comprador escutando em/orders/{id}/events). - O notify_seller vai pra fila — não bloqueia a resposta do checkout.
Idempotência evita decremento duplo de estoque
Se o comprador retentar com a mesma Idempotency-Key (reload, timeout de rede, double-tap), o middleware devolve a resposta original — o handler não roda 2x, então o estoque não é decrementado 2x e nenhum pedido duplicado é criado.
5. Expedição + atualização em tempo real¶
sequenceDiagram
autonumber
actor A as Admin (vendedor)
actor B as Comprador
participant R as orders.router
participant OS as OrderService
participant SS as StockService
participant SSE as orders/{id}/events
participant DB as Postgres
B->>SSE: GET /orders/{id}/events<br/>Accept: text/event-stream
SSE-->>B: event: status (PAID)
Note over A: vendedor expede
A->>R: POST /orders/{id}/ship {tracking}
R->>OS: transition(order_id, SHIPPED)
OS->>OS: assert current == PAID
OS->>DB: UPDATE orders SET status=SHIPPED
OS->>SSE: publish {status: SHIPPED, tracking}
SSE-->>B: event: status (SHIPPED)
Note over A: cliente confirma recebimento
B->>R: POST /orders/{id}/confirm-delivery
R->>OS: transition(order_id, DELIVERED)
OS->>OS: assert current == SHIPPED
OS->>SS: convert_reservation_to_out(items)
SS->>DB: INSERT stock_movements (kind=OUT) por item
OS->>DB: UPDATE orders SET status=DELIVERED
OS->>SSE: publish {status: DELIVERED}
SSE-->>B: event: status (DELIVERED)
Pontos do SDK:
SSEBrokermantém um canal por usuário — cada cliente conectado do comprador recebe o frame (ver fluxo 6 pro fan-out completo).- Transição MUST validar o estado origem (state machine no service).
- Estoque vira
OUTdefinitivo só na entrega — se cancelar antes, oRESERVATIONviraRELEASE.
6. Notificações: SSE + Web Push (um evento, dois canais)¶
Todo evento de domínio relevante pro usuário — pedido pago, pedido expedido, convite recebido, novo review — é entregue em dois canais que carregam o mesmo payload: SSE (SSEBroker, canal = id do usuário) pros clientes com o app aberto (foreground, ao vivo) e Web Push (VAPID) pros dispositivos com o app/aba fechado (background). É "notificação como mensageria": um único NotificationService.notify(...) faz o fan-out pros dois.
sequenceDiagram
autonumber
actor A as Admin (vendedor)
participant OS as OrderService
participant N as NotificationService
participant SB as SSEBroker
participant WP as WebPushSubscriptionService
actor FG as App aberto (foreground)
actor BG as Dispositivo fechado (background)
Note over A: evento de domínio: pedido expedido
A->>OS: POST /orders/{id}/ship {tracking}
OS->>OS: transition(order_id, SHIPPED)
OS->>N: notify(buyer_id, "order_shipped", title, body, data)
par canal foreground
N->>SB: publish(str(user_id), data, event="order_shipped")
SB-->>FG: frame SSE (event: order_shipped)
and canal background
N->>WP: notify_user(user_id, WebPushPayloadSchema(...))
WP-->>BG: Web Push (VAPID) → service worker exibe a notificação
end
Note over WP: dispositivos mortos (404/410) são podados automaticamente
O que "fan-out" quer dizer aqui: o mesmo evento de domínio sai por
dois canais diferentes carregando o payload idêntico — SSE pro app aberto
(entrega ao vivo, agora) e Web Push pro dispositivo fechado (o service worker
acorda e mostra a notificação). Quem chama notify(...) não escolhe o canal:
os dois disparam sempre, com o mesmo data. Assim o app aberto e o app fechado
recebem exatamente a mesma informação — só muda como ela aparece.
Passo 1 — o produtor dispara um único evento¶
O produtor é um controller/service. Logo depois da transição de domínio (aqui,
o pedido virou SHIPPED), ele faz uma só chamada — não sabe nem se importa
que existem dois canais por baixo:
from tempest_fastapi_sdk import BaseService
from src.db.models import OrderModel
from src.db.repositories import OrderRepository
from src.schemas import OrderResponseSchema
from src.services.notifications import NotificationService
class OrderService(BaseService[OrderRepository, OrderResponseSchema]):
def __init__(
self, repository: OrderRepository, notifications: NotificationService
) -> None:
"""Guarda o repositório e o serviço de notificação."""
super().__init__(repository)
self.notifications = notifications
async def mark_shipped(self, order: OrderModel) -> None:
"""Depois da transição de domínio, avisa o comprador — uma chamada só."""
await self.notifications.notify(
order.buyer_id,
event="order_shipped",
title="Pedido a caminho",
body=f"Pedido {order.code} saiu para entrega.",
data={"order_id": str(order.id), "status": order.status},
)
O que cada argumento faz:
order.buyer_id— quem recebe a notificação. Esse id vira o canal SSE e o alvo do Web Push. É obuyer_id(o comprador), não o vendedor que acabou de expedir: quem precisa saber que o pedido saiu é o comprador. É o mesmo id que o comprador usa pra assinar o stream (Passo 3) — os dois lados precisam combinar na mesma string de canal.event="order_shipped"— o nome do evento. No SSE vira oevent:do frame (o front escuta comaddEventListener("order_shipped", ...)); no Web Push vira atagdo payload.title/body— o texto visível da notificação Web Push, que o service worker mostra com o app fechado. O SSE ao vivo ignora esses campos (o app aberto renderiza a UI dele a partir dodata).data— o payload de máquina, igual nos dois canais. É o que o frontend lê pra saber qual pedido mudou e pra qual estado.
Passo 2 — o NotificationService faz o fan-out¶
O NotificationService é o único ponto que conhece os dois canais. Ele
recebe aquela chamada única e a desdobra em dois envios:
# src/services/notification.py
from uuid import UUID
from tempest_fastapi_sdk import SSEBroker, WebPushPayloadSchema, WebPushSubscriptionService
class NotificationService:
"""Fan one domain event out to SSE (foreground) and Web Push (background)."""
def __init__(self, broker: SSEBroker, push: WebPushSubscriptionService) -> None:
self.broker = broker
self.push = push
async def notify(
self, user_id: UUID, event: str, title: str, body: str, data: dict
) -> None:
"""Deliver one event on both channels with the same payload."""
await self.broker.publish(str(user_id), data, event=event)
await self.push.notify_user(
user_id, WebPushPayloadSchema(title=title, body=body, tag=event, data=data)
)
As duas linhas de notify(), uma de cada vez:
await self.broker.publish(str(user_id), data, event=event)— o canal foreground (ao vivo). Mandadatapra todos os streams SSE inscritos no canalstr(user_id)— ou seja, cada aba/app que o usuário tem aberto naquele instante. É fire-and-forget: se ninguém está conectado, opublishnão faz nada (não dá erro). Repare que o canal é a versão string do mesmouser_id.await self.push.notify_user(user_id, WebPushPayloadSchema(...))— o canal background. Envia um Web Push (VAPID) pra todos os dispositivos que o usuário registrou; o service worker de cada um exibe a notificação mesmo com o app/aba fechado. Dispositivos mortos (respondem 404/410) são podados automaticamente do banco. OWebPushPayloadSchemaembrulhatitle/body(texto visível),tag=event(agrupa notificações do mesmo tipo) e o mesmodataque foi pro SSE.
O ponto-chave: o data é o mesmo objeto nos dois envios. O app aberto o
recebe via SSE e o app fechado via Web Push, mas o conteúdo é idêntico — por isso
o frontend consegue tratar os dois com um handler só.
Passo 3 — o cliente assina o stream (GET /notifications/stream)¶
Do outro lado, o cliente com o app aberto assina o canal dele para receber a
metade SSE. O endpoint é uma linha: broker.response(str(user.id)).
# src/api/routers/notifications.py
from fastapi import APIRouter, Depends
from starlette.responses import StreamingResponse
from tempest_fastapi_sdk import SSEBroker
from src.api.dependencies.auth import get_current_user
from src.api.dependencies.resources import get_broker
from src.db.models import User
router = APIRouter()
@router.get("/notifications/stream")
async def notifications_stream(
user: User = Depends(get_current_user),
broker: SSEBroker = Depends(get_broker),
) -> StreamingResponse:
"""Subscribe the caller to their own notification channel and stream it."""
return broker.response(str(user.id))
O que acontece a cada GET /notifications/stream:
Depends(get_current_user)resolve quem está pedindo o stream. O id dele é a chave de tudo — é o mesmo id que o produtor usou como canal no Passo 1.Depends(get_broker)entrega o mesmoSSEBrokersingleton que oNotificationServiceusa pra publicar (sem isso, opublishe oresponsefalariam com brokers diferentes e nada chegaria).broker.response(str(user.id))faz três coisas numa chamada só:- register — cria um
EventStreamnovo e o inscreve no canaluser.id; - stream — devolve uma
StreamingResponsecom os headers de SSE já prontos, e o cliente começa a receber os frames; - unregister — amarra um
on_disconnectque remove esse stream do canal quando o cliente cai, então não sobra registro vazando.
- register — cria um
A partir daí, todo broker.publish(str(user.id), ...) do Passo 2 pinga neste
stream. Um frame SSE que chega nele:
event: order_shipped
id: 01J8Z9F2K7Q3M5R8T0W1X2Y3Z4
data: {"order_id": "9f8e7d6c-5b4a-3210-fedc-ba9876543210", "status": "SHIPPED"}
Pontos do SDK:
- SSE é core (sem extra):
SSEBroker(),await broker.publish(channel, data, event=..., id=..., retry=...),broker.response(channel)já monta aStreamingResponseque assina o canal e desregistra ao desconectar. SSE multi-worker precisa deSSEBroker(redis=...)+broker.run()no lifespan → extra[cache]. - Web Push precisa do extra
[webpush](uv add "tempest-fastapi-sdk[webpush]"): monteWebPushDispatcher(**settings.webpush_kwargs()), passe-o proWebPushSubscriptionService(repository, dispatcher);await service.notify_user(user_id, payload, *, ttl_seconds=None, exclude_endpoints=None)envia pra todos os dispositivos do usuário e poda os mortos (404/410). - O mesmo
dataviaja nos dois canais — o frontend trata SSE e Web Push com o mesmo handler. - Detalhes dos primitivos: Receita SSE » e Receita Web Push ».
Máquina de estados — Order¶
stateDiagram-v2
[*] --> PENDING : checkout
PENDING --> PAID : payment confirmed (admin mock)
PENDING --> CANCELLED : buyer/admin cancels
PAID --> SHIPPED : seller ships
PAID --> CANCELLED : refund pre-ship
SHIPPED --> DELIVERED : buyer confirms
SHIPPED --> RETURNED : return flow
DELIVERED --> [*]
CANCELLED --> [*]
RETURNED --> [*]
Transições inválidas devem falhar com ConflictException
Transições proibidas (qualquer outra setinha) MUST falhar com ConflictException("invalid state transition"). Implementação típica num enum + dict[from, set[to]] no service.
Máquina de estados — Invitation¶
stateDiagram-v2
[*] --> PENDING : invited
PENDING --> ACCEPTED : invitee accepts
PENDING --> REVOKED : admin revokes
PENDING --> EXPIRED : job 7d
PENDING --> SUPERSEDED : new invite for same email
ACCEPTED --> [*]
REVOKED --> [*]
EXPIRED --> [*]
SUPERSEDED --> [*]
EXPIRED é set por tarefa TaskIQ que roda de hora em hora varrendo convites com expires_at < now().
Próximo passo¶
Pula pro Mapa de endpoints ver a API REST completa pronta pra cabear contratos no frontend.