Skip to content

Commit 80551ec

Browse files
committed
fix(service): surface duplicate dispatch conflicts as Task errors
1 parent 755653e commit 80551ec

3 files changed

Lines changed: 105 additions & 61 deletions

File tree

‎core/service/Service/Application.hs‎

Lines changed: 37 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -1026,21 +1026,13 @@ runWithResolved eventStore maybeFileUploadSetup fileUploadCleanup maybeWebAuthSe
10261026
wiredQueries
10271027
|> Array.reduce (\(_reg, (name, _handler, schema)) acc -> acc |> Map.set name schema) Map.empty
10281028

1029-
-- 4. Create query subscriber with combined registry
1030-
subscriber <- Subscriber.new eventStore combinedRegistry
1031-
1032-
-- 5. Rebuild all queries from historical events
1033-
Subscriber.rebuildAll subscriber
1034-
1035-
-- 6. Start live subscription
1036-
Subscriber.start subscriber
1037-
1038-
-- 7. Collect command endpoints and schemas from all services, grouped by transport
1029+
-- 4. Collect command endpoints and schemas from all services, grouped by transport
10391030
endpointsAndSchemasByTransport <-
10401031
app.serviceRunners
10411032
|> Task.mapArray (\runner -> runner.getEndpointsByTransport eventStore app.transports)
10421033

1043-
-- 8. Merge all public endpoints and schemas by transport name, plus dispatch handlers
1034+
-- 5. Merge all public endpoints and schemas by transport name, plus dispatch handlers
1035+
-- Duplicate dispatch handler names are surfaced as Task errors (not panics).
10441036
let mergeEndpointsAndSchemas (serviceEndpoints, serviceSchemas, serviceDispatch) (handlersAcc, schemasAcc, dispatchAcc) = do
10451037
let mergedHandlers = serviceEndpoints
10461038
|> Map.entries
@@ -1060,20 +1052,45 @@ runWithResolved eventStore maybeFileUploadSetup fileUploadCleanup maybeWebAuthSe
10601052
Just existing -> innerAcc |> Map.set transportName (Map.merge cmdSchemas existing)
10611053
)
10621054
schemasAcc
1063-
let mergedDispatch = serviceDispatch
1055+
let mergedDispatchResult = serviceDispatch
10641056
|> Map.entries
10651057
|> Array.reduce
1066-
( \(commandName, handler) innerAcc ->
1067-
case innerAcc |> Map.get commandName of
1068-
Nothing -> innerAcc |> Map.set commandName handler
1069-
Just _ -> panic [fmt|Duplicate command handler registered for integration dispatch: #{commandName}|]
1058+
( \(commandName, handler) innerResult ->
1059+
case innerResult of
1060+
Err err -> Err err
1061+
Ok innerAcc ->
1062+
case innerAcc |> Map.get commandName of
1063+
Nothing -> Ok (innerAcc |> Map.set commandName handler)
1064+
Just _ -> Err [fmt|Duplicate command handler registered for integration dispatch: #{commandName}|]
10701065
)
1071-
dispatchAcc
1072-
(mergedHandlers, mergedSchemas, mergedDispatch)
1066+
(Ok dispatchAcc)
1067+
case mergedDispatchResult of
1068+
Err err -> Err err
1069+
Ok mergedDispatch -> Ok (mergedHandlers, mergedSchemas, mergedDispatch)
10731070

1074-
let (combinedEndpointsByTransport, combinedSchemasByTransport, combinedCommandEndpoints) =
1071+
let mergedEndpointsResult =
10751072
endpointsAndSchemasByTransport
1076-
|> Array.reduce mergeEndpointsAndSchemas (Map.empty, Map.empty, Map.empty)
1073+
|> Array.reduce
1074+
(\serviceResult accResult ->
1075+
case accResult of
1076+
Err err -> Err err
1077+
Ok acc -> mergeEndpointsAndSchemas serviceResult acc
1078+
)
1079+
(Ok (Map.empty, Map.empty, Map.empty))
1080+
1081+
(combinedEndpointsByTransport, combinedSchemasByTransport, combinedCommandEndpoints) <-
1082+
case mergedEndpointsResult of
1083+
Err err -> Task.throw err
1084+
Ok merged -> Task.yield merged
1085+
1086+
-- 6. Create query subscriber with combined registry
1087+
subscriber <- Subscriber.new eventStore combinedRegistry
1088+
1089+
-- 7. Rebuild all queries from historical events
1090+
Subscriber.rebuildAll subscriber
1091+
1092+
-- 8. Start live subscription
1093+
Subscriber.start subscriber
10771094

10781095
-- 9. Initialize auth if configured (using pre-resolved WebAuthSetup)
10791096
maybeAuthEnabled <- case maybeWebAuthSetup of

‎core/service/Service/ServiceDefinition/Core.hs‎

Lines changed: 62 additions & 41 deletions
Original file line numberDiff line numberDiff line change
@@ -210,6 +210,52 @@ buildDispatchHandler _ handler requestContext body respond = do
210210
respond (errorResponse, responseJson)
211211

212212

213+
buildCommandResponseHandler ::
214+
forall cmd entity event entityName name.
215+
( Command cmd,
216+
name ~ NameOf cmd,
217+
Record.KnownSymbol name,
218+
Record.KnownHash name,
219+
Entity entity,
220+
Event event,
221+
event ~ EventOf entity,
222+
entity ~ EntityOf cmd,
223+
entity ~ EntityOf event,
224+
IsMultiTenant cmd ~ False,
225+
StreamId.ToStreamId (EntityIdType entity),
226+
Eq (EntityIdType entity),
227+
Ord (EntityIdType entity),
228+
Show (EntityIdType entity),
229+
Show event,
230+
Record.KnownSymbol entityName,
231+
Json.FromJSON event,
232+
Json.ToJSON event
233+
) =>
234+
EventStore event ->
235+
Maybe (SnapshotCache entity) ->
236+
RequestContext ->
237+
cmd ->
238+
Task Text Response.CommandResponse
239+
buildCommandResponseHandler eventStore maybeCache requestContext cmdInstance = do
240+
fetcher <- case maybeCache of
241+
Just cache ->
242+
EntityFetcher.newWithCache
243+
eventStore
244+
cache
245+
(initialStateImpl @entity)
246+
(updateImpl @entity)
247+
|> Task.mapError toText
248+
Nothing ->
249+
EntityFetcher.new
250+
eventStore
251+
(initialStateImpl @entity)
252+
(updateImpl @entity)
253+
|> Task.mapError toText
254+
let entityName = EntityName (getSymbolText (Record.Proxy @entityName))
255+
result <- CommandExecutor.execute eventStore fetcher entityName requestContext cmdInstance
256+
Task.yield (Response.fromExecutionResult result)
257+
258+
213259
instance
214260
( Command cmd,
215261
Entity entity,
@@ -260,48 +306,16 @@ instance
260306
createHandlers _ eventStore maybeCache transportsMap cmd = do
261307
-- Build the handler that will receive RequestContext at call time
262308
let handler :: RequestContext -> cmd -> Task Text Response.CommandResponse
263-
handler requestContext cmdInstance = do
264-
fetcher <- case maybeCache of
265-
Just cache ->
266-
EntityFetcher.newWithCache
267-
eventStore
268-
cache
269-
(initialStateImpl @entity)
270-
(updateImpl @entity)
271-
|> Task.mapError toText
272-
Nothing ->
273-
EntityFetcher.new
274-
eventStore
275-
(initialStateImpl @entity)
276-
(updateImpl @entity)
277-
|> Task.mapError toText
278-
let entityName = EntityName (getSymbolText (Record.Proxy @(entityName)))
279-
result <- CommandExecutor.execute eventStore fetcher entityName requestContext cmdInstance
280-
Task.yield (Response.fromExecutionResult result)
309+
handler requestContext cmdInstance =
310+
buildCommandResponseHandler @cmd @entity @event @entityName eventStore maybeCache requestContext cmdInstance
281311

282312
buildHandlersForAll @(PublicTransports transports) transportsMap cmd handler
283313

284314

285315
createDispatchHandler _ eventStore maybeCache cmd = do
286316
let handler :: RequestContext -> cmd -> Task Text Response.CommandResponse
287-
handler requestContext cmdInstance = do
288-
fetcher <- case maybeCache of
289-
Just cache ->
290-
EntityFetcher.newWithCache
291-
eventStore
292-
cache
293-
(initialStateImpl @entity)
294-
(updateImpl @entity)
295-
|> Task.mapError toText
296-
Nothing ->
297-
EntityFetcher.new
298-
eventStore
299-
(initialStateImpl @entity)
300-
(updateImpl @entity)
301-
|> Task.mapError toText
302-
let entityName = EntityName (getSymbolText (Record.Proxy @(entityName)))
303-
result <- CommandExecutor.execute eventStore fetcher entityName requestContext cmdInstance
304-
Task.yield (Response.fromExecutionResult result)
317+
handler requestContext cmdInstance =
318+
buildCommandResponseHandler @cmd @entity @event @entityName eventStore maybeCache requestContext cmdInstance
305319

306320
buildDispatchHandler @cmd cmd handler
307321

@@ -565,15 +579,22 @@ buildEndpointsByTransport rawEventStore commandDefinitions transportsMap = do
565579
let newSchemas = schemasAcc |> Map.set transportNameText (existingSchemas |> Map.set cmdName endpointSchema)
566580
(newHandlers, newSchemas)
567581

568-
let groupDispatch (cmdName, handler) dispatchAcc =
569-
case dispatchAcc |> Map.get cmdName of
570-
Just _ -> panic [fmt|Duplicate command handler registered for integration dispatch: #{cmdName}|]
571-
Nothing -> dispatchAcc |> Map.set cmdName handler
582+
let groupDispatch (cmdName, handler) dispatchResult =
583+
case dispatchResult of
584+
Err err -> Err err
585+
Ok dispatchAcc ->
586+
case dispatchAcc |> Map.get cmdName of
587+
Just _ -> Err [fmt|Duplicate command handler registered for integration dispatch: #{cmdName}|]
588+
Nothing -> Ok (dispatchAcc |> Map.set cmdName handler)
572589

573590
let (handlersByTransport, schemasByTransport) =
574591
endpoints |> Array.reduce groupByTransport (Map.empty, Map.empty)
575592

576-
let dispatchMap = dispatchHandlers |> Array.reduce groupDispatch Map.empty
593+
let dispatchMapResult = dispatchHandlers |> Array.reduce groupDispatch (Ok Map.empty)
594+
595+
dispatchMap <- case dispatchMapResult of
596+
Err err -> Task.throw err
597+
Ok mapValue -> Task.yield mapValue
577598

578599
Task.yield (handlersByTransport, schemasByTransport, dispatchMap)
579600

‎docs/decisions/0050-internal-command-transport.md‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -122,6 +122,10 @@ module Service.Transport.Internal (
122122
data InternalTransport = InternalTransport
123123

124124
type instance NameOf InternalTransport = "InternalTransport"
125+
```
126+
127+
```haskell
128+
-- Runtime hooks are owned by service runtime modules, not Service.Transport.Internal.
125129

126130
dispatchCommand ::
127131
Map Text EndpointHandler ->
@@ -135,6 +139,8 @@ buildEndpointsByTransport ::
135139
Task Text (Map Text (Map Text EndpointHandler), Map Text (Map Text EndpointSchema), Map Text EndpointHandler)
136140
```
137141

142+
`dispatchCommand` is implemented in `Service.Integration.Dispatcher` and `buildEndpointsByTransport` is implemented in `Service.ServiceDefinition.Core`.
143+
138144
Implementation consequences of this API:
139145

140146
- `Service.ServiceDefinition.Core` must filter public transports away from `InternalTransport` while still building dispatch handlers for commands that use it.

0 commit comments

Comments
 (0)