Skip to content

Commit 27a68b3

Browse files
Bill LeoutsakosBill Leoutsakos
authored andcommitted
fix(oci-streaming): normalize absent pool settings and filters
1 parent ce20cd1 commit 27a68b3

4 files changed

Lines changed: 133 additions & 6 deletions

File tree

‎apps/sim/blocks/blocks/oci_streaming.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -882,7 +882,7 @@ export const OciStreamingBlock: BlockConfig<OciStreamingResponse> = {
882882
'ociRegion',
883883
'requestId',
884884
]) {
885-
result[key] = params[key] === '' ? undefined : params[key]
885+
result[key] = params[key] === '' || params[key] == null ? undefined : params[key]
886886
}
887887
result.partitions = parseOptionalNumberInput(params.partitions, 'partitions')
888888
result.retentionInHours = parseOptionalNumberInput(

‎apps/sim/lib/internal/oci-streaming/execute-tool.test.ts‎

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,9 @@ vi.mock('@/lib/internal/oci/client.server', () => ({ createOciClient: mocks.crea
1515
import { AuthType } from '@/lib/auth/hybrid'
1616
import { executeOciStreamingTool } from '@/lib/internal/oci-streaming/execute-tool'
1717
import type { InternalToolOperationCall } from '@/lib/internal/tool-operations/types'
18+
import { OciStreamingBlock } from '@/blocks/blocks/oci_streaming'
19+
import { ociStreamingListStreamPoolsTool } from '@/tools/oci_streaming/list_stream_pools'
20+
import type { OciStreamingListStreamPoolsParams } from '@/tools/oci_streaming/types'
1821

1922
function call(overrides: Partial<InternalToolOperationCall> = {}): InternalToolOperationCall {
2023
return {
@@ -52,6 +55,38 @@ describe('OCI Streaming trusted credential boundary', () => {
5255
})
5356
})
5457

58+
it.each([null, ''])(
59+
'omits blank workflow filters after merging the block patch: %s',
60+
async (blank) => {
61+
const inputs = {
62+
operation: 'oci_streaming_list_stream_pools',
63+
ociCredential: 'supplied-reference',
64+
compartmentId: 'compartment-1',
65+
name: blank,
66+
page: blank,
67+
lifecycleState: blank,
68+
sortBy: blank,
69+
sortOrder: blank,
70+
}
71+
const transformed = OciStreamingBlock.tools.config!.params!(inputs)
72+
const operationInput = ociStreamingListStreamPoolsTool.operation.input({
73+
...inputs,
74+
...transformed,
75+
} as OciStreamingListStreamPoolsParams)
76+
const response = await executeOciStreamingTool(
77+
call({ toolId: 'oci_streaming_list_stream_pools', input: operationInput })
78+
)
79+
expect(response.status).toBe(200)
80+
expect(mocks.request).toHaveBeenCalledTimes(1)
81+
expect(mocks.request.mock.calls[0][0].queryPairs).not.toEqual(
82+
expect.arrayContaining([['name', 'null']])
83+
)
84+
expect(mocks.request.mock.calls[0][0].queryPairs).not.toEqual(
85+
expect.arrayContaining([['page', '']])
86+
)
87+
}
88+
)
89+
5590
it('passes only the resolved credential ID and trusted workspace to the foundation', async () => {
5691
const response = await executeOciStreamingTool(call())
5792
expect(response.status).toBe(200)

‎apps/sim/lib/internal/oci-streaming/operations.test.ts‎

Lines changed: 75 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,81 @@ function harness() {
6464
afterEach(() => vi.useRealTimers())
6565

6666
describe('OCI Streaming request and response semantics', () => {
67+
it.each(['create_stream_pool', 'get_stream_pool', 'update_stream_pool'])(
68+
'projects inactive pool settings for %s without replaying the request',
69+
async (operation) => {
70+
const h = harness()
71+
const pool = {
72+
id: 'pool',
73+
name: 'events',
74+
compartmentId: 'compartment',
75+
lifecycleState: 'ACTIVE',
76+
timeCreated: '2026-01-01T00:00:00Z',
77+
kafkaSettings: { bootstrapServers: null },
78+
customEncryptionKey: { kmsKeyId: null, keyState: 'NONE' },
79+
privateEndpointSettings: { nsgIds: null, privateEndpointIp: null, subnetId: null },
80+
}
81+
h.request.mockResolvedValue(response(JSON.stringify(pool), { etag: 'v1' }))
82+
const params =
83+
operation === 'create_stream_pool'
84+
? { name: 'events', compartmentId: 'compartment' }
85+
: operation === 'update_stream_pool'
86+
? { streamPoolId: 'pool', name: 'events' }
87+
: { streamPoolId: 'pool' }
88+
const result = await h.execute({ operation, ...params })
89+
expect(result.output).toMatchObject({
90+
streamPool: {
91+
...pool,
92+
kafkaSettings: { bootstrapServers: undefined },
93+
customEncryptionKey: { kmsKeyId: undefined, keyState: 'NONE' },
94+
privateEndpointSettings: {
95+
nsgIds: undefined,
96+
privateEndpointIp: undefined,
97+
subnetId: undefined,
98+
},
99+
},
100+
etag: 'v1',
101+
})
102+
expect(h.request).toHaveBeenCalledTimes(1)
103+
}
104+
)
105+
106+
it('preserves populated pool settings and rejects malformed settings', async () => {
107+
const h = harness()
108+
const pool = {
109+
id: 'pool',
110+
name: 'events',
111+
compartmentId: 'compartment',
112+
lifecycleState: 'ACTIVE',
113+
timeCreated: '2026-01-01T00:00:00Z',
114+
kafkaSettings: { bootstrapServers: 'server:9092' },
115+
customEncryptionKey: { kmsKeyId: 'key' },
116+
privateEndpointSettings: {
117+
nsgIds: ['nsg'],
118+
privateEndpointIp: '10.0.0.1',
119+
subnetId: 'subnet',
120+
},
121+
}
122+
h.request.mockResolvedValue(response(JSON.stringify(pool)))
123+
expect(
124+
(await h.execute({ operation: 'get_stream_pool', streamPoolId: 'pool' })).output.streamPool
125+
).toEqual(pool)
126+
h.request.mockResolvedValue(
127+
response(JSON.stringify({ ...pool, privateEndpointSettings: { nsgIds: 'nsg' } }))
128+
)
129+
await expect(
130+
h.execute({ operation: 'get_stream_pool', streamPoolId: 'pool' })
131+
).rejects.toThrow()
132+
expect(
133+
ociStreamingInputSchema.safeParse({
134+
operation: 'create_stream_pool',
135+
ociCredential: 'credential',
136+
name: 'events',
137+
compartmentId: 'compartment',
138+
customEncryptionKeyDetails: { kmsKeyId: null },
139+
}).success
140+
).toBe(false)
141+
})
67142
it('accepts full resource names, bounds group paths, and preserves empty tag updates', () => {
68143
const create = {
69144
operation: 'create_stream',

‎apps/sim/lib/internal/oci-streaming/schema.ts‎

Lines changed: 22 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -392,16 +392,33 @@ export const streamPoolSummarySchema = z.object({ ...resource, isPrivate: z.bool
392392
export const streamPoolSchema = streamPoolSummarySchema.extend({
393393
endpointFqdn: z.string().nullish(),
394394
lifecycleStateDetails: z.string().nullish(),
395-
kafkaSettings: kafkaSettings.strip().extend({ bootstrapServers: z.string().optional() }),
395+
kafkaSettings: kafkaSettings.strip().extend({
396+
bootstrapServers: z
397+
.string()
398+
.nullish()
399+
.transform((value) => value ?? undefined),
400+
}),
396401
customEncryptionKey: z.object({
397-
kmsKeyId: z.string().optional(),
402+
kmsKeyId: z
403+
.string()
404+
.nullish()
405+
.transform((value) => value ?? undefined),
398406
keyState: z.string().optional(),
399407
}),
400408
privateEndpointSettings: z
401409
.object({
402-
nsgIds: z.array(z.string()).optional(),
403-
privateEndpointIp: z.string().optional(),
404-
subnetId: z.string().optional(),
410+
nsgIds: z
411+
.array(z.string())
412+
.nullish()
413+
.transform((value) => value ?? undefined),
414+
privateEndpointIp: z
415+
.string()
416+
.nullish()
417+
.transform((value) => value ?? undefined),
418+
subnetId: z
419+
.string()
420+
.nullish()
421+
.transform((value) => value ?? undefined),
405422
})
406423
.nullish(),
407424
})

0 commit comments

Comments
 (0)