u70b3 commented on code in PR #68781:
URL: https://github.com/apache/doris/pull/68781#discussion_r4229151953


##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceIndexAdmission.java:
##########
@@ -141,70 +112,43 @@ static Outcome admitCreate(SnapshotLoader loader, 
LanceExternalCatalog catalog,
         }
         String storedName = LanceIndexFamilies.uniqueMatch(storedNames, 
normalizedName);
         // 4. IF preflight (design section 2.2).
-        boolean orReplace = def.isOrReplace();
-        if (!orReplace && storedName != null) {
+        if (!def.isOrReplace() && storedName != null) {
             if (!ifNotExists) {
                 rejectInvalid("index '" + displayName + "' already exists");
             }
             if (!matchesExistingDefinition(snapshot, storedName, def)) {
                 rejectInvalid("index '" + displayName + "' already exists with 
a different definition");
             }
-            return catalogMgr.withLanceIndexAdmission(catalog, target, () -> 
new Outcome(null));
+            catalogMgr.withLanceIndexAdmission(catalog, target, () -> null);
+            return;
         }
         // 5. Fail closed when the table lookup relation cannot resolve the 
request column
         // uniquely: a dataset can hold top-level fields that differ only by 
case (V versus v),
         // and ExternalTable.getColumn returns the first hit, so admitting 
would journal an
         // arbitrary one of them. The pinned snapshot decides, never the 
cached table schema.
         rejectIfAmbiguousLookupColumn(snapshot, def.getCols().get(0));
-        // 6. Schema contract v1 from the stored column name (never the raw 
user spelling).
-        String storedColumnName = storedColumnName(table, 
def.getCols().get(0));
-        LanceIndexSchemaContract contract =
-                LanceSchemaContractBuilder.build(snapshot.getTopLevelFields(), 
storedColumnName);
-        // 7. The fence locator is the normalized dataset uri of the same 
pinned snapshot.
-        String locator = normalizeLocator(snapshot);
-        // 8. Deterministic normalized properties JSON for ANN; scalar 
families persist null.
-        boolean ann = def.getLanceIndexType() == null;
-        String indexType = ann ? annIndexType(def) : def.getLanceIndexType();
-        String propertiesJson = ann ? buildAnnPropertiesJson(def) : null;
-        // D7 backstop: quota values from fe.conf bypass the ADMIN SET 
callback, so admission
-        // re-asserts positivity before any id allocation or durable transfer.
-        assertPositiveQuotas();
-        return catalogMgr.withLanceIndexAdmission(catalog, target, () -> {
-            // 9. Exactly one id allocation, after every preflight above has 
passed.
-            long jobId = Env.getCurrentEnv().getNextId();
-            String creator = ConnectContext.get().getQualifiedUser();
-            // 10. REPLACE on an existing name persists the stored display 
name (section 4.1) so the
-            // worker locates the case-sensitive target; a fresh REPLACE keeps 
the user's spelling.
-            String persistedDisplayName = (orReplace && storedName != null) ? 
storedName : displayName;
-            LanceIndexJob job;
-            try {
-                job = new LanceIndexJob(jobId, creator, catalog.getId(), 
db.getFullName(), table.getName(),
-                        LanceIndexFenceKey.PROVIDER_DIRECTORY, locator, 
persistedDisplayName, normalizedName,
-                        orReplace ? LanceIndexJobMutationType.REPLACE : 
LanceIndexJobMutationType.CREATE,
-                        ifNotExists, false, indexType, storedColumnName, 
propertiesJson,
-                        snapshot.getDatasetVersion(), contract);
-            } catch (IllegalArgumentException e) {
-                throw invalidAdmission(e.getMessage());
-            }
-            Env.getCurrentEnv().getLanceIndexJobManager().createJob(job,
-                    Config.lance_index_job_max_unresolved_per_table,
-                    Config.lance_index_job_max_unresolved_per_catalog,
-                    Config.lance_index_job_max_unresolved_global);
-            // 11. The job and its fence are durable once createJob returns.
-            return new Outcome(jobId);
+        // 6. Schema contract v1 from the stored column name (never the raw 
user spelling). The
+        // build itself is the admission-depth type validation against the 
pinned snapshot.
+        LanceSchemaContractBuilder.build(snapshot.getTopLevelFields(), 
storedColumnName(table, def.getCols().get(0)));
+        // Every validation above is preserved for the synchronous execution 
path; until that
+        // path lands, a mutation that would be admitted terminates with the 
shared rejection.
+        catalogMgr.withLanceIndexAdmission(catalog, target, () -> {
+            LanceIndexMutationValidator.rejectUnsupportedOperation(

Review Comment:
   Fixed in db84fc9. `DorisFlightSqlProducer.executeQueryStatementLocked` now 
checks the forwarded proxy status after `handleQuery`: a statement the master 
rejected sets the local state from the master error code/message and fails 
GetFlightInfo with INTERNAL before the OK result is synthesized. Covers direct 
and prepared execution; a forwarded statement the master accepted still returns 
the OK StatusResult. New case in DorisFlightSqlSchemaTest (direct Execute, 
prepared Execute with the handle surviving, accepted forward).



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to