Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -145,3 +145,68 @@ func TestSpecWithNoSurfaceGetsNoManifest(t *testing.T) {
t.Error("ops.ts was emitted for a spec with neither endpoints nor channels")
}
}

// A channel that declares a send and a receive operation and no entity is a
// duplex channel: the client speaks on it and takes the frames raw. The
// generator derives that from the AsyncAPI alone, so a backend that already
// declares both operations needs no extension key to reach the client.
func TestDuplexChannelGetsADuplexBinding(t *testing.T) {
spec := streamOnlySpec()
spec.WebSockets = append(spec.WebSockets, client.WebSocketEndpoint{
ID: "liveQueryWS",
Path: "/api/v1/query/live/ws",
SendSchema: &client.Schema{Type: "object"},
ReceiveSchema: &client.Schema{Type: "object"},
Metadata: map[string]any{
"messages": map[string]string{"send": "send", "receive": "receive"},
},
})

config := baseConfig()
config.Hooks = true
config.IncludeStreaming = true

out, err := NewGenerator().Generate(context.Background(), spec, config)
if err != nil {
t.Fatalf("Generate: %v", err)
}

ops := ClientManifestText(out.Files)

for _, want := range []string{
"kind: 'duplex'",
"channel: '/api/v1/query/live/ws'",
"send: 'send'",
"receive: 'receive'",
"kind: 'entity'",
} {
if !strings.Contains(ops, want) {
t.Errorf("stream bindings are missing %q in:\n%s", want, ops)
}
}
}

// A receive-only channel without an entity is neither: it stays out of the
// table exactly as before, so a listen-only endpoint does not become a
// subscribable channel by accident.
func TestReceiveOnlyChannelWithoutEntityStaysOut(t *testing.T) {
spec := streamOnlySpec()
spec.WebSockets = append(spec.WebSockets, client.WebSocketEndpoint{
ID: "ticker",
Path: "/api/v1/ticker",
ReceiveSchema: &client.Schema{Type: "object"},
Metadata: map[string]any{"messages": map[string]string{"tick": "receive"}},
})

config := baseConfig()
config.IncludeStreaming = true

out, err := NewGenerator().Generate(context.Background(), spec, config)
if err != nil {
t.Fatalf("Generate: %v", err)
}

if strings.Contains(ClientManifestText(out.Files), "/api/v1/ticker") {
t.Errorf("a receive-only channel without an entity must not be emitted")
}
}
59 changes: 57 additions & 2 deletions internal/client/generators/typescript/opsmanifest.go
Original file line number Diff line number Diff line change
Expand Up @@ -680,12 +680,30 @@ func (g *OpsManifestGenerator) writeStreams(buf *strings.Builder, spec *client.A
bindings []client.StreamBinding
}

type duplex struct {
path string
send string
receive string
}

channels := make([]channel, 0, len(spec.WebSockets)+len(spec.SSEs))
duplexes := make([]duplex, 0)

for i := range spec.WebSockets {
if b := spec.WebSockets[i].StreamBindings; len(b) > 0 {
channels = append(channels, channel{spec.WebSockets[i].Path, b})
ws := &spec.WebSockets[i]

if b := ws.StreamBindings; len(b) > 0 {
channels = append(channels, channel{ws.Path, b})

continue
}

if ws.SendSchema == nil || ws.ReceiveSchema == nil {
continue
}

send, receive := duplexMessageNames(ws.Metadata)
duplexes = append(duplexes, duplex{ws.Path, send, receive})
}

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

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

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

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

// A duplex channel: the client speaks on it and no entity stands behind it.
// Derived from the declared operations alone, so the backend needs no
// extension key for a channel it already describes as send and receive.
for _, d := range duplexes {
buf.WriteString(" {\n")
buf.WriteString(" kind: 'duplex',\n")
buf.WriteString(fmt.Sprintf(" channel: %s,\n", tsString(d.path)))
buf.WriteString(fmt.Sprintf(" send: %s,\n", tsString(d.send)))
buf.WriteString(fmt.Sprintf(" receive: %s,\n", tsString(d.receive)))
buf.WriteString(" },\n")
}

buf.WriteString("] as const;\n")
}

// duplexMessageNames picks the lowest-sorted message name for each direction,
// the same tie-break convertSchemaFromChannel uses, so the output is stable
// across runs.
func duplexMessageNames(metadata map[string]any) (string, string) {
names, _ := metadata["messages"].(map[string]string)
send, receive := "", ""

for _, name := range sortedKeys(names) {
switch names[name] {
case "send":
if send == "" {
send = name
}
case "receive":
if receive == "" {
receive = name
}
}
}

return send, receive
}

// tsString renders a single-quoted TypeScript string literal.
//
// Escaping is not paranoia: a path or tag reaching the file unescaped closes the
Expand Down
80 changes: 80 additions & 0 deletions packages/client-core/__tests__/duplex-binding.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,80 @@
import { describe, expect, it } from 'vitest';

import { QueryCache } from '../src/cache';
import { manualScheduler } from '../src/invalidate';
import { binderSnapshot, StreamBinder } from '../src/live';
import { SubscriptionManager } from '../src/stream';
import type { StreamBinding } from '../src/stream';
import { manualClock } from '../src/transport';
import { fakeSockets, fakeTransport } from './harness';

const streams: readonly StreamBinding[] = [
{ kind: 'duplex', channel: '/api/v1/query/live/ws', send: 'SendMessage', receive: 'ReceiveMessage' },
{ channel: '/ws/orders', message: 'orderUpdated', entity: 'Order', intent: 'upsert', invalidates: [] },
];

function build() {
const sockets = fakeSockets();
const clock = manualClock();
const cache = new QueryCache({ transport: fakeTransport(() => []), entities: { Order: { idField: 'id' } } });
const manager = new SubscriptionManager({
connect: sockets.connect,
sleep: clock.sleep,
random: () => 0,
release: manualScheduler().schedule,
});
const binder = new StreamBinder({ cache, streams, manager, scheduler: (flush) => flush(), sleep: clock.sleep });

return { binder, sockets, manager };
}

describe('a duplex channel', () => {
// Guards that a raw subscription is the manager's subscription: same socket,
// same hello, same release.
it('subscribes raw with a hello and hands frames over undecoded', () => {
const { binder, sockets } = build();
const seen: unknown[] = [];

const release = binder.raw('/api/v1/query/live/ws', (message) => seen.push(message), {
hello: { action: 'subscribe', data: { id: 'q1' } },
goodbye: { action: 'unsubscribe', data: { id: 'q1' } },
});
sockets.last().open();
sockets.last().deliver({ type: 'snapshot', subscriptionId: 'q1', payload: { rows: [] } });

expect(sockets.last().sent).toEqual([{ action: 'subscribe', data: { id: 'q1' } }]);
expect(seen).toEqual([{ type: 'snapshot', subscriptionId: 'q1', payload: { rows: [] } }]);

release();
expect(sockets.last().sent).toEqual([
{ action: 'subscribe', data: { id: 'q1' } },
{ action: 'unsubscribe', data: { id: 'q1' } },
]);
});

// Guards the generated table as the single source of channel names.
it('refuses a channel that is not a duplex binding', () => {
const { binder } = build();

expect(() => binder.raw('/ws/orders', () => {})).toThrow(/not a duplex channel/);
expect(() => binder.raw('/ws/nowhere', () => {})).toThrow(/not a duplex channel/);
});

// Guards the entity path from a frame it must never see.
it('keeps duplex frames out of the entity store', () => {
const { binder, sockets } = build();

binder.raw('/api/v1/query/live/ws', () => {}, { hello: { id: 'q1' } });
sockets.last().open();
sockets.last().deliver({ type: 'orderUpdated', payload: { id: 'o1', total: 3 } });

expect(binderSnapshot(binder).queued).toBe(0);
});

it('lists the duplex binding in the snapshot', () => {
const { binder } = build();

const channels = binderSnapshot(binder).channels.map((c) => c.channel).sort();
expect(channels).toEqual(['/api/v1/query/live/ws', '/ws/orders']);
});
});
Loading
Loading