Repository navigation
Expand file tree
/
Copy pathlocal.ts
More file actions
156 lines (144 loc) · 6.05 KB
/
Copy pathlocal.ts
File metadata and controls
156 lines (144 loc) · 6.05 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
// Shared client for the local Maple server's `POST /local/query` endpoint, used
// by both the browser SPA (`apps/local-ui`) and the query CLI (`apps/cli`).
// The endpoint runs raw SQL through the in-process chDB session and returns a
// bare JSON array.
//
// The output FORMAT is owned by the server: `apps/cli/src/server/query-guard.ts`
// peels whatever trailing `FORMAT <fmt>` the compiler emitted (`CH.compile(...)`
// appends `FORMAT JSON`) and runs the single read-only statement as JSON rows.
// So callers POST `compiled.sql` verbatim.
import { Effect, Schema } from "effect"
import { FetchHttpClient, HttpClient, HttpClientRequest } from "effect/http"
import { OrgId } from "@maple/domain/primitives"
/**
* Local mode's single synthetic tenant. The ingest binary writes every span,
* log and metric under this `OrgId`, so every local query compiles with it.
*/
export const LOCAL_ORG_ID = OrgId.make("local")
/**
* The local server refused the query. Its own tag, and structured fields, so
* callers stop re-deriving the cause from the rendered sentence: `status` is
* the HTTP status, `detail` the server's body, and `code`/`type` the chDB error
* identity lifted out of it (`60` / `UNKNOWN_TABLE`).
*/
export class LocalQueryFailed extends Schema.TaggedError<LocalQueryFailed>()(
"@maple/query-engine/LocalQueryFailed",
{
status: Schema.Number,
detail: Schema.String,
code: Schema.optionalKey(Schema.String),
type: Schema.optionalKey(Schema.String),
message: Schema.String,
},
) {}
/** The server answered 2xx with something that was not the documented JSON array. */
export class LocalQueryMalformedResponse extends Schema.TaggedError<LocalQueryMalformedResponse>()(
"@maple/query-engine/LocalQueryMalformedResponse",
{ message: Schema.String, cause: Schema.optionalKey(Schema.Defect()) },
) {}
/** The request never reached a response: the binary is down, or the socket dropped. */
export class LocalQueryUnreachable extends Schema.TaggedError<LocalQueryUnreachable>()(
"@maple/query-engine/LocalQueryUnreachable",
{ message: Schema.String, cause: Schema.Defect() },
) {}
export type LocalQueryError = LocalQueryFailed | LocalQueryMalformedResponse | LocalQueryUnreachable
export type LocalQueryRow = Record<string, unknown>
const LocalQueryRows = Schema.Array(Schema.Record(Schema.String, Schema.Unknown))
const decodeRows = Schema.decodeUnknownEffect(LocalQueryRows)
/**
* Execute compiled SQL against the local Maple binary and return the rows.
*
* Lazy and interruptible: the request is bound to a scope that closes with
* success, failure, or interruption, so a cancelled caller aborts the HTTP
* request instead of leaving chDB working on a result nobody reads.
*
* @param sql The compiled SQL (e.g. from `CH.compile(...).sql`), sent as-is.
* @param baseUrl Origin of the local binary. Defaults to `""` (a relative
* `/local/query`, for the SPA behind its vite proxy); the CLI
* passes an absolute address like `http://127.0.0.1:4318`.
*/
export const executeLocalQuery = (
sql: string,
baseUrl = "",
): Effect.Effect<ReadonlyArray<LocalQueryRow>, LocalQueryError, HttpClient.HttpClient> =>
Effect.scoped(
Effect.gen(function* () {
const http = yield* HttpClient.HttpClient
const request = HttpClientRequest.post(`${baseUrl}/local/query`).pipe(
HttpClientRequest.bodyText(JSON.stringify({ sql }), "application/json"),
)
const response = yield* HttpClient.withScope(http)
.execute(request)
.pipe(
Effect.mapError(
(cause) =>
new LocalQueryUnreachable({
message: "Local Maple server is unreachable",
cause,
}),
),
)
if (response.status < 200 || response.status >= 300) {
const detail = (yield* response.text.pipe(Effect.orElseSucceed(() => ""))).trim()
return yield* new LocalQueryFailed({
status: response.status,
detail,
// `code`/`type` are what `mapWarehouseError` classifies on. Lifting them
// here lets the classifier see `UNKNOWN_TABLE` instead of regex-matching
// the rendered sentence.
...clickHouseErrorFields(detail),
message: `Local query failed (${response.status})${detail ? `: ${detail}` : ""}`,
})
}
const json = yield* response.json.pipe(
Effect.mapError(
(cause) =>
new LocalQueryMalformedResponse({
message: "Local query response was not JSON",
cause,
}),
),
)
return yield* decodeRows(json).pipe(
Effect.mapError(
(cause) =>
new LocalQueryMalformedResponse({
message: "Local query response was not a JSON array of rows",
cause,
}),
),
)
}),
)
/**
* Promise edge for the SPA's hooks, which hand TanStack Query an `AbortSignal`.
* Runs `executeLocalQuery` on the platform `fetch`; aborting the signal
* interrupts the fiber, which aborts the request.
*/
export const runLocalQuery = (
sql: string,
baseUrl = "",
signal?: AbortSignal,
): Promise<ReadonlyArray<LocalQueryRow>> =>
// This is the SPA's entry point into Effect: each call is its own root, so
// the fetch layer is provided here rather than composed higher up.
// oxlint-disable-next-line effecttsgo/strict-effect-provide
Effect.runPromise(executeLocalQuery(sql, baseUrl).pipe(Effect.provide(FetchHttpClient.layer)), { signal })
/**
* chDB renders its failures as `query failed: Code: 60. DB::Exception: … (UNKNOWN_TABLE)`.
* Lift the numeric code and the symbolic type out of that text so the error
* carries them as fields.
*/
type ChdbErrorIdentity = { code?: string; type?: string }
const clickHouseErrorFields = (detail: string): ChdbErrorIdentity => {
// Assigned rather than spread conditionally: both are `optionalKey`, so an
// explicit `undefined` is not the same as an absent key.
const fields: ChdbErrorIdentity = {}
const code = detail.match(/\bCode:\s*(\d+)/)?.[1]
if (code !== undefined) fields.code = code
// The type is the trailing parenthesised SCREAMING_CASE token; chDB puts it
// last, after the human sentence.
const type = detail.match(/\(([A-Z][A-Z0-9_]{2,})\)\s*$/)?.[1]
if (type !== undefined) fields.type = type
return fields
}