Skip to content

Commit 58dda85

Browse files
SK-3026 add tests
1 parent 5a4a8f5 commit 58dda85

13 files changed

Lines changed: 570 additions & 16 deletions

File tree

‎common/src/main/java/com/skyflow/BaseSkyflow.java‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -69,7 +69,10 @@ protected static <T> T resolveOrThrow(Map<String, T> map, String key,
6969
ErrorLogs errorLog, ErrorMessage errorMessage) throws SkyflowException {
7070
T value = key != null ? map.get(key) : map.values().stream().findFirst().orElse(null);
7171
if (value == null) {
72-
LogUtil.printErrorLog(errorLog.getLog());
72+
// The log line carries a %s1 placeholder for the id. Callers that resolve the single
73+
// configured entry pass no key, so say so rather than emitting the raw placeholder.
74+
LogUtil.printErrorLog(BaseUtils.parameterizedString(
75+
errorLog.getLog(), key != null ? key : "not specified"));
7376
throw new SkyflowException(ErrorCode.INVALID_INPUT.getCode(), errorMessage.getMessage());
7477
}
7578
return value;

‎common/src/main/java/com/skyflow/vault/data/RequestContext.java‎

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,15 +7,36 @@
77
import java.util.Map;
88

99
public final class RequestContext {
10+
/** Reported when the caller's request was not split into batches. */
11+
private static final int NOT_BATCHED = -1;
12+
1013
private final String operation;
14+
private final int batchIndex;
15+
private final int totalBatches;
1116
private final Map<CustomHeaderKey, String> headers = new HashMap<>();
1217

1318
public RequestContext(String operation) {
19+
this(operation, NOT_BATCHED, NOT_BATCHED);
20+
}
21+
22+
public RequestContext(String operation, int batchIndex, int totalBatches) {
1423
this.operation = operation;
24+
this.batchIndex = batchIndex;
25+
this.totalBatches = totalBatches;
1526
}
1627

1728
public String getOperation() { return operation; }
1829

30+
/**
31+
* Zero-based position of this batch within the caller's request, or -1 when the operation was
32+
* not batched. Lets an interceptor tag each batch distinctly — a per-batch correlation id, for
33+
* instance — instead of seeing an identical context for every one.
34+
*/
35+
public int getBatchIndex() { return batchIndex; }
36+
37+
/** Total number of batches the request was split into, or -1 when it was not batched. */
38+
public int getTotalBatches() { return totalBatches; }
39+
1940
public void addHeader(CustomHeaderKey key, String value) {
2041
headers.put(key, value);
2142
}

‎flowvault/src/main/java/com/skyflow/Skyflow.java‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,21 @@ public VaultConfig getVaultConfig() {
4040
return (VaultConfig) array[0];
4141
}
4242

43+
/**
44+
* Updates a vault's configuration on an already-built client.
45+
* <p>
46+
* BaseSkyflow.updateVaultConfig goes straight to the template, bypassing the builder's own
47+
* override, so the flowvault-specific fields have to be carried across here too — otherwise
48+
* a vaultUrl or HTTP setting supplied through this entry point would be silently dropped
49+
* while the same call on the builder honoured it.
50+
*/
51+
@Override
52+
public Skyflow updateVaultConfig(VaultConfig vaultConfig) throws SkyflowException {
53+
super.updateVaultConfig(vaultConfig);
54+
this.builder.carryVaultOverrides(vaultConfig);
55+
return this;
56+
}
57+
4358
public VaultController vault() throws SkyflowException {
4459
return resolveOrThrow(this.builder.vaultClientsMap, null, ErrorLogs.VAULT_CONFIG_DOES_NOT_EXIST, ErrorMessage.VaultIdNotInConfigList);
4560
}

‎flowvault/src/main/java/com/skyflow/utils/Utils.java‎

Lines changed: 8 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -636,18 +636,15 @@ private static String tokenAt(List<String> tokens, int position) {
636636

637637
private static BulkDeleteTokensResponseRecord createDeleteTokensErrorRecord(
638638
Map<String, Object> recordMap, int index, String requestedToken, String requestId) {
639-
int code = 500;
640-
if (recordMap.containsKey("http_code")) {
641-
code = (Integer) recordMap.get("http_code");
642-
} else if (recordMap.containsKey("httpCode")) {
643-
code = (Integer) recordMap.get("httpCode");
644-
} else if (recordMap.containsKey("statusCode")) {
645-
code = (Integer) recordMap.get("statusCode");
639+
// Read through the shared helpers rather than casting: recordMap holds deserialised JSON,
640+
// so a status can arrive as Double or String depending on the parser, and a blind
641+
// (Integer) cast would turn a real API error into a ClassCastException.
642+
int code = readHttpCode(recordMap, 500);
643+
String message = readErrorMessage(recordMap);
644+
String token = readString(recordMap, "value");
645+
if (token == null) {
646+
token = requestedToken;
646647
}
647-
String message = recordMap.containsKey("error") ? (String) recordMap.get("error") :
648-
recordMap.containsKey("message") ? (String) recordMap.get("message") : "Unknown error";
649-
Object echoedToken = recordMap.get("value");
650-
String token = (echoedToken instanceof String) ? (String) echoedToken : requestedToken;
651648
return new BulkDeleteTokensResponseRecord(index, token, code, message, requestId);
652649
}
653650

‎flowvault/src/main/java/com/skyflow/vault/controller/VaultController.java‎

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -460,7 +460,7 @@ private List<CompletableFuture<BulkDeleteTokensResponse>> deleteTokensBatchFutur
460460
for (int batchIndex = 0; batchIndex < batches.size(); batchIndex++) {
461461
final int index = batchIndex;
462462
com.skyflow.generated.rest.resources.flowservice.requests.V1FlowDeleteTokenRequest batch = batches.get(index);
463-
RequestContext ctx = new RequestContext("DELETE_TOKENS");
463+
RequestContext ctx = new RequestContext("DELETE_TOKENS", batchIndex, batches.size());
464464
if (interceptor != null) interceptor.intercept(ctx);
465465
CompletableFuture<BulkDeleteTokensResponse> future = CompletableFuture
466466
.supplyAsync(() -> processDeleteTokensBatch(batch, ctx), executor)
@@ -612,12 +612,14 @@ private List<CompletableFuture<BulkTokenizeResponse>> tokenizeBatchFutures(
612612
// batches are contiguous but not uniformly sized - a batch is cut short when it would
613613
// otherwise repeat a value - so track where each one starts rather than deriving it
614614
int nextStartIndex = 0;
615+
int batchPosition = 0;
615616
for (List<BulkTokenizeRequestRecord> batchRecords : batches) {
616617
final int startIndex = nextStartIndex;
617618
nextStartIndex += batchRecords.size();
619+
final int batchIndex = batchPosition++;
618620
com.skyflow.generated.rest.resources.flowservice.requests.V1FlowTokenizeRequest batch =
619621
Utils.getBulkTokenizeRequestBody(batchRecords, this.getVaultConfig().getVaultId());
620-
RequestContext ctx = new RequestContext("TOKENIZE");
622+
RequestContext ctx = new RequestContext("TOKENIZE", batchIndex, batches.size());
621623
if (interceptor != null) interceptor.intercept(ctx);
622624
CompletableFuture<BulkTokenizeResponse> future = CompletableFuture
623625
.supplyAsync(() -> processTokenizeBatch(batch, ctx), executor)
@@ -830,7 +832,7 @@ private List<CompletableFuture<BulkDetokenizeResponse>> detokenizeBatchFutures(
830832
for (int batchIndex = 0; batchIndex < batches.size(); batchIndex++) {
831833
com.skyflow.generated.rest.resources.flowservice.requests.V1FlowDetokenizeRequest batch = batches.get(batchIndex);
832834
int batchNumber = batchIndex;
833-
RequestContext ctx = new RequestContext("DETOKENIZE");
835+
RequestContext ctx = new RequestContext("DETOKENIZE", batchIndex, batches.size());
834836
if (interceptor != null) interceptor.intercept(ctx);
835837
CompletableFuture<BulkDetokenizeResponse> future = CompletableFuture
836838
.supplyAsync(() -> processDetokenizeBatch(batch, ctx), executor)
@@ -862,7 +864,7 @@ private List<CompletableFuture<BulkInsertResponse>> insertBatchFutures(
862864
for (int batchIndex = 0; batchIndex < batches.size(); batchIndex++) {
863865
List<V1InsertRecordData> batch = batches.get(batchIndex);
864866
int batchNumber = batchIndex;
865-
RequestContext ctx = new RequestContext("INSERT");
867+
RequestContext ctx = new RequestContext("INSERT", batchIndex, batches.size());
866868
if (interceptor != null) interceptor.intercept(ctx);
867869
CompletableFuture<BulkInsertResponse> future = CompletableFuture
868870
.supplyAsync(() -> insertBatch(
Lines changed: 212 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,212 @@
1+
package com.skyflow;
2+
3+
import com.skyflow.config.VaultConfig;
4+
import com.skyflow.enums.Env;
5+
import com.skyflow.errors.SkyflowException;
6+
import com.skyflow.vault.controller.VaultController;
7+
import org.junit.Assert;
8+
import org.junit.Test;
9+
10+
/**
11+
* The four client lifecycle scenarios, mirrored by the ClientOperationsExample sample.
12+
*
13+
* <p>Each scenario runs twice: once through {@code SkyflowClientBuilder} and once through the built
14+
* {@code Skyflow} client. The two go down different code paths — BaseSkyflow.updateVaultConfig calls
15+
* the template directly and skips the builder's own override — and a bug that dropped the
16+
* flowvault-specific fields on the client path only was found exactly this way.
17+
*
18+
* <p>Unlike the sample, these run inside the com.skyflow package, so they can assert against the
19+
* controller's own config rather than only the stored copy.
20+
*/
21+
public class ClientLifecycleScenarioTests {
22+
23+
private static final String VAULT_ID = "vault1";
24+
private static final String NOT_IN_CONFIG_LIST = "VaultId is missing from the config";
25+
26+
private static VaultConfig config(String clusterId) {
27+
VaultConfig config = new VaultConfig();
28+
config.setVaultId(VAULT_ID);
29+
config.setClusterId(clusterId);
30+
config.setEnv(Env.DEV);
31+
return config;
32+
}
33+
34+
/** An update carrying a new cluster, env and timeout. clusterId is resent because the incoming
35+
* config is validated on its own before being merged. */
36+
private static VaultConfig update(String clusterId, Env env, Integer timeout) {
37+
VaultConfig update = config(clusterId);
38+
update.setEnv(env);
39+
update.setTimeout(timeout);
40+
return update;
41+
}
42+
43+
// ── A: add -> update -> delete -> vault() must fail ───────────────────────
44+
45+
@Test
46+
public void testScenarioA_viaBuilder_addUpdateDeleteThenVaultFails() throws SkyflowException {
47+
Skyflow.SkyflowClientBuilder builder = Skyflow.builder().addVaultConfig(config("cluster1"));
48+
49+
// add
50+
Assert.assertEquals("cluster1", builder.build().getVaultConfig(VAULT_ID).getClusterId());
51+
// update - no error
52+
builder.updateVaultConfig(update("cluster2", Env.PROD, 30));
53+
Assert.assertEquals("cluster2", builder.build().getVaultConfig(VAULT_ID).getClusterId());
54+
Assert.assertEquals(Env.PROD, builder.build().getVaultConfig(VAULT_ID).getEnv());
55+
Assert.assertEquals(Integer.valueOf(30), builder.build().getVaultConfig(VAULT_ID).getTimeout());
56+
// delete
57+
builder.removeVaultConfig(VAULT_ID);
58+
Assert.assertNull(builder.build().getVaultConfig(VAULT_ID));
59+
60+
// vault() -> vault id not found
61+
Skyflow client = builder.build();
62+
try {
63+
client.vault();
64+
Assert.fail("vault() must fail once the vault is removed");
65+
} catch (SkyflowException e) {
66+
Assert.assertTrue(e.getMessage().contains(NOT_IN_CONFIG_LIST));
67+
}
68+
}
69+
70+
@Test
71+
public void testScenarioA_viaClient_addUpdateDeleteThenVaultFails() throws SkyflowException {
72+
Skyflow client = Skyflow.builder().addVaultConfig(config("cluster1")).build();
73+
74+
client.updateVaultConfig(update("cluster2", Env.PROD, 30));
75+
Assert.assertEquals("cluster2", client.getVaultConfig(VAULT_ID).getClusterId());
76+
Assert.assertEquals(Integer.valueOf(30), client.getVaultConfig(VAULT_ID).getTimeout());
77+
78+
client.removeVaultConfig(VAULT_ID);
79+
Assert.assertNull(client.getVaultConfig(VAULT_ID));
80+
81+
try {
82+
client.vault();
83+
Assert.fail("vault() must fail once the vault is removed");
84+
} catch (SkyflowException e) {
85+
Assert.assertTrue(e.getMessage().contains(NOT_IN_CONFIG_LIST));
86+
}
87+
}
88+
89+
// ── B: add -> delete -> update must throw ─────────────────────────────────
90+
91+
@Test
92+
public void testScenarioB_viaBuilder_addDeleteThenUpdateThrows() throws SkyflowException {
93+
Skyflow.SkyflowClientBuilder builder = Skyflow.builder()
94+
.addVaultConfig(config("cluster1"))
95+
.removeVaultConfig(VAULT_ID);
96+
97+
try {
98+
builder.updateVaultConfig(update("cluster2", Env.PROD, 30));
99+
Assert.fail("updating a removed vault must throw, not silently re-create it");
100+
} catch (SkyflowException e) {
101+
Assert.assertTrue(e.getMessage().contains(NOT_IN_CONFIG_LIST));
102+
}
103+
104+
// the failed update must not have resurrected the vault
105+
Assert.assertNull(builder.build().getVaultConfig(VAULT_ID));
106+
}
107+
108+
@Test
109+
public void testScenarioB_viaClient_addDeleteThenUpdateThrows() throws SkyflowException {
110+
Skyflow client = Skyflow.builder().addVaultConfig(config("cluster1")).build();
111+
client.removeVaultConfig(VAULT_ID);
112+
113+
try {
114+
client.updateVaultConfig(update("cluster2", Env.PROD, 30));
115+
Assert.fail("updating a removed vault must throw, not silently re-create it");
116+
} catch (SkyflowException e) {
117+
Assert.assertTrue(e.getMessage().contains(NOT_IN_CONFIG_LIST));
118+
}
119+
120+
Assert.assertNull(client.getVaultConfig(VAULT_ID));
121+
try {
122+
client.vault();
123+
Assert.fail("vault() must still fail after the rejected update");
124+
} catch (SkyflowException e) {
125+
Assert.assertTrue(e.getMessage().contains(NOT_IN_CONFIG_LIST));
126+
}
127+
}
128+
129+
// ── C: add -> update -> vault() carries the latest config ─────────────────
130+
131+
@Test
132+
public void testScenarioC_viaBuilder_vaultAfterUpdateHasLatest() throws SkyflowException {
133+
Skyflow.SkyflowClientBuilder builder = Skyflow.builder().addVaultConfig(config("cluster1"));
134+
135+
builder.updateVaultConfig(update("cluster2", Env.PROD, 15));
136+
VaultController vault = builder.build().vault();
137+
138+
Assert.assertEquals("cluster2", vault.getVaultConfig().getClusterId());
139+
Assert.assertEquals(Env.PROD, vault.getVaultConfig().getEnv());
140+
Assert.assertEquals(Integer.valueOf(15), vault.getVaultConfig().getTimeout());
141+
Assert.assertEquals("https://cluster2.skyvault.skyflowapis.com", vault.currentVaultURL);
142+
vault.updateExecutorInHTTP();
143+
Assert.assertEquals(15_000, vault.sharedHttpClient.callTimeoutMillis());
144+
}
145+
146+
@Test
147+
public void testScenarioC_viaClient_vaultAfterUpdateHasLatest() throws SkyflowException {
148+
Skyflow client = Skyflow.builder().addVaultConfig(config("cluster1")).build();
149+
150+
client.updateVaultConfig(update("cluster2", Env.PROD, 15));
151+
VaultController vault = client.vault();
152+
153+
Assert.assertEquals("cluster2", vault.getVaultConfig().getClusterId());
154+
Assert.assertEquals(Env.PROD, vault.getVaultConfig().getEnv());
155+
Assert.assertEquals(Integer.valueOf(15), vault.getVaultConfig().getTimeout());
156+
Assert.assertEquals("https://cluster2.skyvault.skyflowapis.com", vault.currentVaultURL);
157+
vault.updateExecutorInHTTP();
158+
Assert.assertEquals(15_000, vault.sharedHttpClient.callTimeoutMillis());
159+
}
160+
161+
// ── D: add -> vault() -> update -> vault() carries the latest config ──────
162+
163+
@Test
164+
public void testScenarioD_viaBuilder_heldControllerSeesTheUpdate() throws SkyflowException {
165+
Skyflow.SkyflowClientBuilder builder = Skyflow.builder().addVaultConfig(config("cluster1"));
166+
VaultController held = builder.build().vault();
167+
Assert.assertEquals("cluster1", held.getVaultConfig().getClusterId());
168+
Assert.assertEquals("https://cluster1.skyvault.skyflowapis.dev", held.currentVaultURL);
169+
170+
builder.updateVaultConfig(update("cluster2", Env.PROD, 45));
171+
VaultController after = builder.build().vault();
172+
173+
Assert.assertSame("the reference taken before the update must still be current", held, after);
174+
Assert.assertEquals("cluster2", held.getVaultConfig().getClusterId());
175+
Assert.assertEquals(Integer.valueOf(45), held.getVaultConfig().getTimeout());
176+
Assert.assertEquals("https://cluster2.skyvault.skyflowapis.com", held.currentVaultURL);
177+
}
178+
179+
@Test
180+
public void testScenarioD_viaClient_heldControllerSeesTheUpdate() throws SkyflowException {
181+
Skyflow client = Skyflow.builder().addVaultConfig(config("cluster1")).build();
182+
VaultController held = client.vault();
183+
Assert.assertEquals("cluster1", held.getVaultConfig().getClusterId());
184+
185+
client.updateVaultConfig(update("cluster2", Env.PROD, 45));
186+
187+
Assert.assertSame("the reference taken before the update must still be current",
188+
held, client.vault());
189+
Assert.assertEquals("cluster2", held.getVaultConfig().getClusterId());
190+
Assert.assertEquals(Integer.valueOf(45), held.getVaultConfig().getTimeout());
191+
Assert.assertEquals("https://cluster2.skyvault.skyflowapis.com", held.currentVaultURL);
192+
held.updateExecutorInHTTP();
193+
Assert.assertEquals(45_000, held.sharedHttpClient.callTimeoutMillis());
194+
}
195+
196+
// ── the vaultUrl variant of C/D, since it resolves differently to clusterId ──
197+
198+
@Test
199+
public void testScenarioD_viaClient_heldControllerSeesANewVaultUrl() throws SkyflowException {
200+
VaultConfig initial = config("cluster1");
201+
initial.setVaultUrl("https://first.example.com");
202+
Skyflow client = Skyflow.builder().addVaultConfig(initial).build();
203+
VaultController held = client.vault();
204+
Assert.assertEquals("https://first.example.com", held.currentVaultURL);
205+
206+
VaultConfig update = config("cluster1");
207+
update.setVaultUrl("https://second.example.com");
208+
client.updateVaultConfig(update);
209+
210+
Assert.assertEquals("https://second.example.com", held.currentVaultURL);
211+
}
212+
}

0 commit comments

Comments
 (0)