Skip to content

Commit d5c970e

Browse files
authored
feat(client): duplex stream channels with a send path (#107)
* feat(client-core): let a subscription speak first, and again after a reconnect subscribe() takes hello and goodbye frames. hello goes out once the transport reports open and after every reconnect, in subscription order; goodbye goes out on release while the socket is up. A transport without send() refuses a hello at subscribe time rather than dropping it. The snapshot says which channels carry a hello, and attempts may be Infinity, which the reconnect loop already honoured. * fix(client-core): send hello once, order onReconnect after greet, guard goodbye-only Three defects from review. First, a first subscribe on a transport with no onOpen greeted the new hello inside open() and then sent it again in subscribe() itself; subscribe() now skips the explicit send when this call is the one that opened the socket. Second, onReconnect fired right after open() returned even on a transport whose open event has not happened yet, so a consumer refetching on reconnect could run before the hello reintroducing the subscription went out; open() now takes the reconnect's channel list and reports it from the ready path, after greet(), never before. Third, the no-send guard only checked hello, so a goodbye-only subscription on a connection without send() returned normally and later dropped its goodbye in silence; the guard now applies whenever either option is present. Updated the gap-recovery and reconnect tests that relied on the old, premature onReconnect timing to call open() on the reconnected fake connection, matching what a real transport requires before the manager can call it ready. * fix(client-core): carry hello/goodbye through a repartition reopen too repartition() had the same two defects reconnect() just had fixed. It called onReconnect right after open(), before the replacement socket's transport had actually reported open, so a consumer acting on onReconnect could run before the replacement was greeted. And it copied channels and refs onto the replacement socket but never frames, so every surviving subscription's hello and goodbye registration was silently dropped on an identity change: no hello went out on the reopened socket and the server never learned the subscription existed again. frames now carries over in the same insertion order, and repartition routes its onReconnect through open()'s report parameter, so it fires from the ready path after greet(), exactly like a drop-triggered reconnect. * feat(client-core): a duplex binding kind the binder hands over raw StreamBinding is now a union: the entity binding everyone generates today, with an optional kind, and a duplex binding for a channel the client speaks on with no entity behind it. The binder keeps duplex channels out of the entity store and exposes raw(channel, handler, options), which subscribes through the manager with hello and goodbye frames and refuses any channel the generated table does not declare. Three existing tests read StreamBinding fields the brief did not list: frame-ordering.test.ts and overlay.test.ts built entity bindings typed as the bare union and passed them to a frame() helper typed the same way, and live.test.ts read .message off a ChannelBindings.bindings array. All three are narrowed to EntityStreamBinding, the first two by retyping the fixtures and the helper, the last with an isDuplex filter, rather than cast. stream binding measured 4040 B against its 3.75 kB budget once raw() and isDuplex landed. Raised to 4.2 kB, which clears the measurement with about the same headroom the query engine's own budget carries. * feat(client-gen): emit a duplex binding for a send-and-receive channel A WebSocket channel that declares both a send and a receive operation and no entity is one the client speaks on; the binder takes its frames raw. writeStreams now emits it with kind 'duplex' and the two message names, and stamps kind 'entity' on the bindings it already emitted. The kind is optional on entity bindings in client-core, so every table generated before this keeps type-checking. * fix(client-devtools): render duplex stream bindings in the streams tab The bindings table read message, entity, intent and invalidates off every binding. With StreamBinding now a union that no longer compiles, and at runtime the first duplex binding threw and blanked the whole tab. The table narrows on kind: an entity row keeps its columns and a duplex row shows the send and receive message names with the entity columns empty. The socket snapshot rows now carry hello, so the live-panel test expects it, and the streams test narrows its binding assertion for the union.
1 parent d96026c commit d5c970e

17 files changed

Lines changed: 731 additions & 50 deletions

File tree

‎internal/client/generators/typescript/e2e_streamonly_manifest_test.go‎

Lines changed: 65 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -145,3 +145,68 @@ func TestSpecWithNoSurfaceGetsNoManifest(t *testing.T) {
145145
t.Error("ops.ts was emitted for a spec with neither endpoints nor channels")
146146
}
147147
}
148+
149+
// A channel that declares a send and a receive operation and no entity is a
150+
// duplex channel: the client speaks on it and takes the frames raw. The
151+
// generator derives that from the AsyncAPI alone, so a backend that already
152+
// declares both operations needs no extension key to reach the client.
153+
func TestDuplexChannelGetsADuplexBinding(t *testing.T) {
154+
spec := streamOnlySpec()
155+
spec.WebSockets = append(spec.WebSockets, client.WebSocketEndpoint{
156+
ID: "liveQueryWS",
157+
Path: "/api/v1/query/live/ws",
158+
SendSchema: &client.Schema{Type: "object"},
159+
ReceiveSchema: &client.Schema{Type: "object"},
160+
Metadata: map[string]any{
161+
"messages": map[string]string{"send": "send", "receive": "receive"},
162+
},
163+
})
164+
165+
config := baseConfig()
166+
config.Hooks = true
167+
config.IncludeStreaming = true
168+
169+
out, err := NewGenerator().Generate(context.Background(), spec, config)
170+
if err != nil {
171+
t.Fatalf("Generate: %v", err)
172+
}
173+
174+
ops := ClientManifestText(out.Files)
175+
176+
for _, want := range []string{
177+
"kind: 'duplex'",
178+
"channel: '/api/v1/query/live/ws'",
179+
"send: 'send'",
180+
"receive: 'receive'",
181+
"kind: 'entity'",
182+
} {
183+
if !strings.Contains(ops, want) {
184+
t.Errorf("stream bindings are missing %q in:\n%s", want, ops)
185+
}
186+
}
187+
}
188+
189+
// A receive-only channel without an entity is neither: it stays out of the
190+
// table exactly as before, so a listen-only endpoint does not become a
191+
// subscribable channel by accident.
192+
func TestReceiveOnlyChannelWithoutEntityStaysOut(t *testing.T) {
193+
spec := streamOnlySpec()
194+
spec.WebSockets = append(spec.WebSockets, client.WebSocketEndpoint{
195+
ID: "ticker",
196+
Path: "/api/v1/ticker",
197+
ReceiveSchema: &client.Schema{Type: "object"},
198+
Metadata: map[string]any{"messages": map[string]string{"tick": "receive"}},
199+
})
200+
201+
config := baseConfig()
202+
config.IncludeStreaming = true
203+
204+
out, err := NewGenerator().Generate(context.Background(), spec, config)
205+
if err != nil {
206+
t.Fatalf("Generate: %v", err)
207+
}
208+
209+
if strings.Contains(ClientManifestText(out.Files), "/api/v1/ticker") {
210+
t.Errorf("a receive-only channel without an entity must not be emitted")
211+
}
212+
}

‎internal/client/generators/typescript/opsmanifest.go‎

Lines changed: 57 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -680,12 +680,30 @@ func (g *OpsManifestGenerator) writeStreams(buf *strings.Builder, spec *client.A
680680
bindings []client.StreamBinding
681681
}
682682

683+
type duplex struct {
684+
path string
685+
send string
686+
receive string
687+
}
688+
683689
channels := make([]channel, 0, len(spec.WebSockets)+len(spec.SSEs))
690+
duplexes := make([]duplex, 0)
684691

685692
for i := range spec.WebSockets {
686-
if b := spec.WebSockets[i].StreamBindings; len(b) > 0 {
687-
channels = append(channels, channel{spec.WebSockets[i].Path, b})
693+
ws := &spec.WebSockets[i]
694+
695+
if b := ws.StreamBindings; len(b) > 0 {
696+
channels = append(channels, channel{ws.Path, b})
697+
698+
continue
688699
}
700+
701+
if ws.SendSchema == nil || ws.ReceiveSchema == nil {
702+
continue
703+
}
704+
705+
send, receive := duplexMessageNames(ws.Metadata)
706+
duplexes = append(duplexes, duplex{ws.Path, send, receive})
689707
}
690708

691709
for i := range spec.SSEs {
@@ -695,12 +713,14 @@ func (g *OpsManifestGenerator) writeStreams(buf *strings.Builder, spec *client.A
695713
}
696714

697715
sort.Slice(channels, func(i, j int) bool { return channels[i].path < channels[j].path })
716+
sort.Slice(duplexes, func(i, j int) bool { return duplexes[i].path < duplexes[j].path })
698717

699718
buf.WriteString("export const streams = [\n")
700719

701720
for _, ch := range channels {
702721
for _, b := range ch.bindings {
703722
buf.WriteString(" {\n")
723+
buf.WriteString(" kind: 'entity',\n")
704724
buf.WriteString(fmt.Sprintf(" channel: %s,\n", tsString(ch.path)))
705725
buf.WriteString(fmt.Sprintf(" message: %s,\n", tsString(b.Message)))
706726
buf.WriteString(fmt.Sprintf(" entity: %s,\n", tsString(b.EntityType)))
@@ -710,9 +730,44 @@ func (g *OpsManifestGenerator) writeStreams(buf *strings.Builder, spec *client.A
710730
}
711731
}
712732

733+
// A duplex channel: the client speaks on it and no entity stands behind it.
734+
// Derived from the declared operations alone, so the backend needs no
735+
// extension key for a channel it already describes as send and receive.
736+
for _, d := range duplexes {
737+
buf.WriteString(" {\n")
738+
buf.WriteString(" kind: 'duplex',\n")
739+
buf.WriteString(fmt.Sprintf(" channel: %s,\n", tsString(d.path)))
740+
buf.WriteString(fmt.Sprintf(" send: %s,\n", tsString(d.send)))
741+
buf.WriteString(fmt.Sprintf(" receive: %s,\n", tsString(d.receive)))
742+
buf.WriteString(" },\n")
743+
}
744+
713745
buf.WriteString("] as const;\n")
714746
}
715747

748+
// duplexMessageNames picks the lowest-sorted message name for each direction,
749+
// the same tie-break convertSchemaFromChannel uses, so the output is stable
750+
// across runs.
751+
func duplexMessageNames(metadata map[string]any) (string, string) {
752+
names, _ := metadata["messages"].(map[string]string)
753+
send, receive := "", ""
754+
755+
for _, name := range sortedKeys(names) {
756+
switch names[name] {
757+
case "send":
758+
if send == "" {
759+
send = name
760+
}
761+
case "receive":
762+
if receive == "" {
763+
receive = name
764+
}
765+
}
766+
}
767+
768+
return send, receive
769+
}
770+
716771
// tsString renders a single-quoted TypeScript string literal.
717772
//
718773
// Escaping is not paranoia: a path or tag reaching the file unescaped closes the
Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,80 @@
1+
import { describe, expect, it } from 'vitest';
2+
3+
import { QueryCache } from '../src/cache';
4+
import { manualScheduler } from '../src/invalidate';
5+
import { binderSnapshot, StreamBinder } from '../src/live';
6+
import { SubscriptionManager } from '../src/stream';
7+
import type { StreamBinding } from '../src/stream';
8+
import { manualClock } from '../src/transport';
9+
import { fakeSockets, fakeTransport } from './harness';
10+
11+
const streams: readonly StreamBinding[] = [
12+
{ kind: 'duplex', channel: '/api/v1/query/live/ws', send: 'SendMessage', receive: 'ReceiveMessage' },
13+
{ channel: '/ws/orders', message: 'orderUpdated', entity: 'Order', intent: 'upsert', invalidates: [] },
14+
];
15+
16+
function build() {
17+
const sockets = fakeSockets();
18+
const clock = manualClock();
19+
const cache = new QueryCache({ transport: fakeTransport(() => []), entities: { Order: { idField: 'id' } } });
20+
const manager = new SubscriptionManager({
21+
connect: sockets.connect,
22+
sleep: clock.sleep,
23+
random: () => 0,
24+
release: manualScheduler().schedule,
25+
});
26+
const binder = new StreamBinder({ cache, streams, manager, scheduler: (flush) => flush(), sleep: clock.sleep });
27+
28+
return { binder, sockets, manager };
29+
}
30+
31+
describe('a duplex channel', () => {
32+
// Guards that a raw subscription is the manager's subscription: same socket,
33+
// same hello, same release.
34+
it('subscribes raw with a hello and hands frames over undecoded', () => {
35+
const { binder, sockets } = build();
36+
const seen: unknown[] = [];
37+
38+
const release = binder.raw('/api/v1/query/live/ws', (message) => seen.push(message), {
39+
hello: { action: 'subscribe', data: { id: 'q1' } },
40+
goodbye: { action: 'unsubscribe', data: { id: 'q1' } },
41+
});
42+
sockets.last().open();
43+
sockets.last().deliver({ type: 'snapshot', subscriptionId: 'q1', payload: { rows: [] } });
44+
45+
expect(sockets.last().sent).toEqual([{ action: 'subscribe', data: { id: 'q1' } }]);
46+
expect(seen).toEqual([{ type: 'snapshot', subscriptionId: 'q1', payload: { rows: [] } }]);
47+
48+
release();
49+
expect(sockets.last().sent).toEqual([
50+
{ action: 'subscribe', data: { id: 'q1' } },
51+
{ action: 'unsubscribe', data: { id: 'q1' } },
52+
]);
53+
});
54+
55+
// Guards the generated table as the single source of channel names.
56+
it('refuses a channel that is not a duplex binding', () => {
57+
const { binder } = build();
58+
59+
expect(() => binder.raw('/ws/orders', () => {})).toThrow(/not a duplex channel/);
60+
expect(() => binder.raw('/ws/nowhere', () => {})).toThrow(/not a duplex channel/);
61+
});
62+
63+
// Guards the entity path from a frame it must never see.
64+
it('keeps duplex frames out of the entity store', () => {
65+
const { binder, sockets } = build();
66+
67+
binder.raw('/api/v1/query/live/ws', () => {}, { hello: { id: 'q1' } });
68+
sockets.last().open();
69+
sockets.last().deliver({ type: 'orderUpdated', payload: { id: 'o1', total: 3 } });
70+
71+
expect(binderSnapshot(binder).queued).toBe(0);
72+
});
73+
74+
it('lists the duplex binding in the snapshot', () => {
75+
const { binder } = build();
76+
77+
const channels = binderSnapshot(binder).channels.map((c) => c.channel).sort();
78+
expect(channels).toEqual(['/api/v1/query/live/ws', '/ws/orders']);
79+
});
80+
});

0 commit comments

Comments
 (0)