diff --git a/Dockerfile b/Dockerfile index 935a36d..c76fa70 100644 --- a/Dockerfile +++ b/Dockerfile @@ -74,7 +74,8 @@ RUN dart pub get COPY . . RUN dart pub get --offline -RUN dart compile exe bin/grpc_server.dart -o /out/noosphere_roast_server +RUN mkdir -p /out \ + && dart compile exe bin/grpc_server.dart -o /out/noosphere_roast_server FROM docker.io/library/debian:bookworm-slim @@ -91,7 +92,7 @@ COPY --from=frosty-build /out/libfrosty_rust.so /app/build/libfrosty_rust.so COPY --from=secp256k1-build /out/libsecp256k1.so /app/build/libsecp256k1.so ENV LD_LIBRARY_PATH="/app/build:/usr/local/lib" -EXPOSE 50051 +EXPOSE 50051 8080 ENTRYPOINT ["/app/noosphere_roast_server", "--config"] -CMD ["/config/server.yaml"] +CMD ["/config/server.yaml", "--rest-address", "0.0.0.0", "--rest-port", "8080"] diff --git a/README.md b/README.md index c027f13..308a5c6 100644 --- a/README.md +++ b/README.md @@ -8,6 +8,34 @@ A server can be run from a given `GrpcConfig` YAML file using `dart run noosphere_roast_server:grpc_server --config your_config_file_here.yaml`. Alternatively a server may be created using the package as a library. +The server emits `info` logs by default. Use `--log-level` to choose one of +`trace`, `debug`, `info`, `warning`, `error`, `fatal`, or `off`: + +```sh +dart run noosphere_roast_server:grpc_server \ + --config your_config_file_here.yaml \ + --log-level debug +``` + +REST/WebSocket can be enabled for browser clients with `--rest-port`: + +```sh +dart run noosphere_roast_server:grpc_server \ + --config your_config_file_here.yaml \ + --rest-port 8080 \ + --rest-allow-origin '*' +``` + +`--rest-allow-origin '*'` is convenient for local testing, but production +deployments should set `--rest-allow-origin` to the exact frontend origin that +will access the REST/WebSocket API, for example `https://app.example.com`. + +Use `--rest-address 0.0.0.0` when the REST/WebSocket listener must be reachable from +outside the process namespace, such as from a container port mapping. The +default REST/WebSocket bind address is `localhost`. + +The REST/WebSocket API shape is documented in [REST_API_SPEC.md](REST_API_SPEC.md). + ## Podman / Docker Build the image from this repository: @@ -16,6 +44,10 @@ Build the image from this repository: podman build -t noosphere-roast-server . ``` +Rebuild the image after changing local source. The Dockerfile copies this local +repository into the image with `COPY . .`, so an old image will not contain +recent CLI, logging, or REST changes. + The Dockerfile builds `libfrosty_rust.so` from `peercoin/frosty` `v3.0.0`, matching the current `frosty` dependency, and `libsecp256k1.so` from `peercoin/secp256k1-coinlib` `v0.7.0`, matching the current `coinlib` @@ -33,10 +65,20 @@ Run the server with a mounted YAML configuration: ```sh podman run --rm \ -p 50051:50051 \ + -p 8080:8080 \ -v "$PWD/config.yaml:/config/server.yaml:ro,Z" \ noosphere-roast-server ``` +The container starts both gRPC and REST/WebSocket by default. gRPC listens on the port +from the YAML config, or `50051` when `port` is omitted. REST/WebSocket listens on +container port `8080`. + +Port mapping syntax is `host_port:container_port`. If the YAML config says +`port: 443`, the gRPC server listens on container port `443`, so map it with +`-p 50051:443` if clients should connect to host port `50051`. If the YAML +config omits `port` or says `port: 50051`, use `-p 50051:50051`. + The `:Z` suffix relabels the mounted config file so Podman can read it on SELinux-enforcing hosts. Use `:z` instead if the same config file must be shared by multiple containers. @@ -46,13 +88,127 @@ To use a different in-container config path, pass it as the command: ```sh podman run --rm \ -p 50051:50051 \ + -p 8080:8080 \ -v "$PWD/config.yaml:/app/config.yaml:ro,Z" \ - noosphere-roast-server /app/config.yaml + noosphere-roast-server \ + /app/config.yaml --rest-address 0.0.0.0 --rest-port 8080 +``` + +### REST/WebSocket With CORS + +For local testing, allow any browser origin and enable debug logs: + +```sh +podman run --rm \ + -p 50051:50051 \ + -p 8080:8080 \ + -v "$PWD/config.yaml:/config/server.yaml:ro,Z" \ + noosphere-roast-server \ + /config/server.yaml \ + --rest-address 0.0.0.0 \ + --rest-port 8080 \ + --rest-allow-origin '*' \ + --log-level debug +``` + +For production, replace `'*'` with the frontend origin that loads the web app: + +```sh +--rest-allow-origin https://app.example.com +``` + +Only one layer should emit CORS headers. If a reverse proxy such as Caddy is +already adding `Access-Control-Allow-Origin`, run the backend with +`--rest-disable-cors` instead. Otherwise browsers will reject responses with a +combined value such as `*, *`. + +### Caddy Reverse Proxy + +Bind container ports to localhost when Caddy runs on the same host: + +```sh +podman run --rm \ + -p 127.0.0.1:50051:50051 \ + -p 127.0.0.1:8080:8080 \ + -v "$PWD/config.yaml:/config/server.yaml:ro,Z" \ + noosphere-roast-server \ + /config/server.yaml \ + --rest-address 0.0.0.0 \ + --rest-port 8080 \ + --rest-allow-origin https://app.example.com \ + --log-level info +``` + +REST/WebSocket on a dedicated API hostname: + +```caddyfile +api.example.com { + reverse_proxy 127.0.0.1:8080 { + flush_interval -1 + } +} +``` + +Do not add CORS headers in both Caddy and the backend. Either let the backend +handle CORS with `--rest-allow-origin`, or let Caddy handle it and run the +backend with `--rest-disable-cors`. + +If the browser frontend is served from the same hostname and REST is under a +prefix, strip the prefix before proxying: + +```caddyfile +app.example.com { + handle_path /api/noosphere/* { + reverse_proxy 127.0.0.1:8080 { + flush_interval -1 + } + } + + root * /srv/app + file_server +} +``` + +### Ngrok For REST/WebSocket Testing + +Expose the REST/WebSocket port, not the gRPC port: + +```sh +ngrok http 8080 +``` + +Use the printed HTTPS URL as the REST base URL in the frontend. The websocket +event stream will be under: + +```text +wss:///sessions//events +``` + +### Logging + +Use `--log-level debug` when diagnosing frontend connectivity: + +```sh +--log-level debug +``` + +At `info`, the server logs lifecycle and coordinator state changes such as +startup, auth challenges, participant login/logout, DKG requests, signature +completion, and shutdown. + +At `debug`, the gRPC transport also logs request receipt/completion and event +stream lifecycle, for example: + +```text +gRPC login received +gRPC login completed +gRPC fetchEventStream opened for participant ... ``` -The image builds the `frosty` and `secp256k1-coinlib` native libraries during -the container build and copies `libfrosty_rust.so` and `libsecp256k1.so` into -`/app/build`. +If shared coordinator logs appear but no `gRPC ... received` logs appear while +running with `--log-level debug`, the frontend is probably using REST or the +gRPC request is not reaching this container. Check the configured client port, +container port mapping, firewall, and any reverse proxy. The same commands also work with Docker by replacing `podman` with `docker`. diff --git a/REST_API_SPEC.md b/REST_API_SPEC.md new file mode 100644 index 0000000..5523c5f --- /dev/null +++ b/REST_API_SPEC.md @@ -0,0 +1,389 @@ +# Noosphere ROAST Server REST/WebSocket API + +This specification is for implementing a frontend REST/WebSocket adapter for the +Noosphere ROAST server. The same frontend may already support the gRPC endpoint; +reuse the same Noosphere domain serializers and parsers where possible. + +## Transport Model + +Base URL: the REST server origin configured separately from gRPC with +`--rest-port`. + +All `POST` endpoints: + +- Request body is JSON. +- Request header should include `Content-Type: application/json`. +- Binary/domain objects are base64 strings of the same `.toBytes()` payloads + used by the gRPC client. +- The REST API takes the same binary payloads gRPC sends as protobuf `bytes`, + then base64-encodes them so JSON can carry them. +- The server accepts standard base64 or URL-safe base64, with or without + padding. +- Do not send gRPC/protobuf wrapper messages to REST; send the underlying + domain bytes encoded as base64. + +Success response shapes: + +```json +{} +``` + +```json +{ "data": "" } +``` + +```json +{ "data": ["", "..."] } +``` + +Error response shapes: + +```json +{ "error": "" } +``` + +Invalid requests return HTTP `400`. Unexpected server errors return HTTP `500` +with `{ "error": "Internal server error" }`. + +## Endpoints + +### POST /login + +Request: + +```json +{ + "groupFingerprint": "", + "participantId": "", + "protocolVersion": 2 +} +``` + +`protocolVersion` is optional and defaults to `2`. + +Response: + +```json +{ "data": "" } +``` + +### POST /respond-to-challenge + +Request: + +```json +{ + "challenge": "", + "signature": "" +} +``` + +Response: + +```json +{ "data": "" } +``` + +Use the decoded `LoginCompleteResponse.id` as the session id for later calls +and the event websocket. + +### POST /extend-session + +Request: + +```json +{ "sid": "" } +``` + +Response: + +```json +{ "data": "" } +``` + +### POST /dkg/new + +Request: + +```json +{ + "sid": "", + "signedDetails": " bytes>", + "commitment": "" +} +``` + +Response: + +```json +{} +``` + +### POST /dkg/reject + +Request: + +```json +{ + "sid": "", + "name": "dkg-name" +} +``` + +Response: + +```json +{} +``` + +### POST /dkg/commitment + +Request: + +```json +{ + "sid": "", + "name": "dkg-name", + "commitment": "" +} +``` + +Response: + +```json +{} +``` + +### POST /dkg/round2 + +Request: + +```json +{ + "sid": "", + "name": "dkg-name", + "commitmentSetSignature": "", + "secrets": [ + { + "id": "", + "secret": "" + } + ] +} +``` + +Response: + +```json +{} +``` + +### POST /dkg/acks + +Request: + +```json +{ + "sid": "", + "acks": [""] +} +``` + +Response: + +```json +{} +``` + +### POST /dkg/request-acks + +Request: + +```json +{ + "sid": "", + "requests": [""] +} +``` + +Response: + +```json +{ "data": [""] } +``` + +### POST /signatures/request + +Request: + +```json +{ + "sid": "", + "keys": [""], + "signedDetails": " bytes>", + "commitments": [""] +} +``` + +Response: + +```json +{} +``` + +### POST /signatures/reject + +Request: + +```json +{ + "sid": "", + "reqId": "" +} +``` + +Response: + +```json +{} +``` + +### POST /signatures/replies + +Request: + +```json +{ + "sid": "", + "reqId": "", + "replies": [""] +} +``` + +Response when new rounds are created: + +```json +{ + "type": "new_round", + "data": "" +} +``` + +Response when signatures are complete: + +```json +{ + "type": "complete", + "data": "" +} +``` + +Response when there is no immediate data: + +```json +{ + "type": "empty", + "data": null +} +``` + +### POST /secret-share + +Request: + +```json +{ + "sid": "", + "groupKey": "", + "secrets": [ + { + "id": "", + "share": "" + } + ] +} +``` + +Response: + +```json +{ "data": [""] } +``` + +### POST /key-constructed/ack + +Request: + +```json +{ + "sid": "", + "constructedKey": " bytes>" +} +``` + +Response: + +```json +{} +``` + +## WebSocket Event Stream + +Open after login: + +```text +GET /sessions//events +``` + +`` is the raw `SessionID.n` bytes encoded as URL-safe base64 with padding +removed. + +Open this endpoint as a websocket. Use `ws://` for plain HTTP deployments and +`wss://` when the REST server is served over HTTPS. + +Each websocket message is a JSON text frame: + +```json +{ + "type": "dkg_commitment", + "data": "" +} +``` + +`type` is the event name. `data` is the matching Noosphere `Event.toBytes()` +payload encoded as base64. + +Event type mapping: + +```text +participant_status +new_dkg +dkg_commitment +dkg_reject +dkg_round2_share +dkg_ack +dkg_ack_request +signatures_request +signature_new_rounds +signatures_complete +signatures_failure +secret_share +constructed_key +keepalive +``` + +## Frontend Adapter Guidance + +Implement a REST adapter with the same method surface as the existing gRPC +adapter. Most methods should be thin wrappers: + +1. Serialize existing Noosphere domain objects to bytes. +2. Base64 encode those bytes into the documented JSON fields. +3. `POST` the JSON request. +4. Decode returned `data` bytes back into the same domain response classes. +5. For websocket events, route by the JSON `type` value and parse `data` as the + matching Event bytes. + +The REST transport does not replace client-side protocol logic. DKG, +signature-round, authentication, and key-sharing behavior should remain the +same as in the gRPC adapter. diff --git a/bin/grpc_server.dart b/bin/grpc_server.dart index fe6c9b1..04e60e2 100644 --- a/bin/grpc_server.dart +++ b/bin/grpc_server.dart @@ -4,6 +4,16 @@ import 'package:args/args.dart'; import 'package:coinlib/coinlib.dart'; import 'package:noosphere_roast_server/noosphere_roast_server.dart'; +const _logLevels = { + "trace": Level.trace, + "debug": Level.debug, + "info": Level.info, + "warning": Level.warning, + "error": Level.error, + "fatal": Level.fatal, + "off": Level.off, +}; + void main(List args) async { final argParser = ArgParser(); argParser.addOption( @@ -12,21 +22,72 @@ void main(List args) async { help: "The path to the GrpcConfig YAML file", mandatory: true, ); + argParser.addOption( + "rest-port", + help: "Optional REST/WebSocket port for browser clients", + ); + argParser.addOption( + "rest-address", + help: "REST/WebSocket bind address", + defaultsTo: "localhost", + ); + argParser.addOption( + "rest-allow-origin", + help: "CORS Access-Control-Allow-Origin value for REST/WebSocket clients", + defaultsTo: "*", + ); + argParser.addFlag( + "rest-disable-cors", + help: "Do not emit CORS headers; use when a reverse proxy handles CORS", + defaultsTo: false, + negatable: false, + ); + argParser.addOption( + "log-level", + help: "Minimum log level to emit", + allowed: _logLevels.keys, + defaultsTo: "info", + ); final argResults = argParser.parse(args); + final logger = createNoosphereRoastServerLogger( + level: _logLevels[argResults.option("log-level")]!, + ); final configFile = argResults.option("config")!; + final restPortString = argResults.option("rest-port"); + final restPort = restPortString == null ? null : int.parse(restPortString); final configString = File(configFile).readAsStringSync(); await loadFrosty(); final config = GrpcConfig.fromYaml(configString); - print("Loaded config from $configFile"); - print("Group fingerprint is ${bytesToHex(config.server.group.fingerprint)}"); + logger.i("Loaded config from $configFile"); + logger.i( + "Group fingerprint is ${bytesToHex(config.server.group.fingerprint)}", + ); - final apiHandler = ServerApiHandler(config: config.server); + final apiHandler = SynchronizedServerApiHandler( + config: config.server, + logger: logger, + ); final service = FrostNoosphereService(api: apiHandler); final grpcServer = service.createServer(); await grpcServer.serve(port: config.port); - print("Server listening on port ${config.port}"); + logger.i("gRPC server listening on port ${config.port}"); + + HttpServer? restServer; + if (restPort != null) { + final restService = RestWebSocketNoosphereService( + api: apiHandler, + allowOrigin: argResults.flag("rest-disable-cors") + ? null + : argResults.option("rest-allow-origin")!, + ); + final restAddress = argResults.option("rest-address")!; + restServer = await restService.serve(address: restAddress, port: restPort); + logger.i( + "REST/WebSocket server listening on $restAddress:${restServer.port}", + ); + } // Wait for SIGINT or SIGTERM to terminate server @@ -35,7 +96,7 @@ void main(List args) async { for (final signal in [ProcessSignal.sigint, ProcessSignal.sigterm]) { signal.watch().listen((sig) { if (termCompleter.isCompleted) { - print("Exiting immediately"); + logger.w("Exiting immediately"); exit(0); } termCompleter.complete(sig); @@ -43,9 +104,12 @@ void main(List args) async { } final signal = await termCompleter.future; - print("Caught ${signal.name}. Shutting down server."); + logger.i( + "Caught ${signal.name}. Shutting down server.", + ); await apiHandler.shutdown(); + await restServer?.close(force: true); await grpcServer.shutdown(); exit(0); diff --git a/lib/src/common.dart b/lib/src/common.dart new file mode 100644 index 0000000..48a09c3 --- /dev/null +++ b/lib/src/common.dart @@ -0,0 +1,7 @@ +import 'dart:typed_data'; +import 'package:noosphere_roast_client/noosphere_roast_client.dart'; + +Uint8List bytes(List li) => Uint8List.fromList(li); +SessionID sid(List li) => SessionID.fromBytes(bytes(li)); +SignaturesRequestId sigReqId(List li) => + SignaturesRequestId.fromBytes(bytes(li)); diff --git a/lib/src/config/grpc.dart b/lib/src/config/grpc.dart index 1fc2a88..7750f83 100644 --- a/lib/src/config/grpc.dart +++ b/lib/src/config/grpc.dart @@ -4,12 +4,14 @@ import 'package:noosphere_roast_client/noosphere_roast_client.dart'; import 'server.dart'; class GrpcConfig with cl.Writable, MapWritable { + static const defaultPort = 50051; + final ServerConfig server; final int port; GrpcConfig({ required this.server, - required this.port, + this.port = defaultPort, }); GrpcConfig.fromReader(cl.BytesReader reader) @@ -28,7 +30,7 @@ class GrpcConfig with cl.Writable, MapWritable { GrpcConfig.fromMapReader(MapReader reader) : this( server: ServerConfig.fromMapReader(reader["server"]), - port: reader["port"].require(), + port: reader["port"].value() ?? defaultPort, ); GrpcConfig.fromYaml(String yaml) diff --git a/lib/src/grpc.dart b/lib/src/grpc.dart index e368283..61a0119 100644 --- a/lib/src/grpc.dart +++ b/lib/src/grpc.dart @@ -1,40 +1,63 @@ import 'dart:async'; -import 'dart:typed_data'; import 'package:coinlib/coinlib.dart' as cl; import 'package:grpc/grpc.dart' as grpc; import 'package:noosphere_roast_client/pbgrpc.dart' as pb; import 'package:noosphere_roast_client/noosphere_roast_client.dart'; +import 'package:noosphere_roast_server/src/common.dart' as common; +import 'package:noosphere_roast_server/src/logging.dart'; import 'package:noosphere_roast_server/src/server/api_handler.dart'; import 'package:noosphere_roast_server/src/server/state/client_session.dart'; - -Uint8List _bytes(List li) => Uint8List.fromList(li); -SessionID _sid(List li) => SessionID.fromBytes(_bytes(li)); -SignaturesRequestId _sigReqId(List li) => - SignaturesRequestId.fromBytes(_bytes(li)); pb.Bytes _returnWritable(cl.Writable writable) => pb.Bytes( data: writable.toBytes(), ); class FrostNoosphereService extends pb.NoosphereServiceBase { final ServerApiHandler api; + final Logger logger; - FrostNoosphereService({required this.api}); + FrostNoosphereService({ + required this.api, + Logger? logger, + }) : logger = logger ?? api.logger; grpc.Server createServer() => grpc.Server.create(services: [this]); - grpc.GrpcError _wrapException(Exception e) => - grpc.GrpcError.unknown(e.toString()); + grpc.GrpcError _wrapException( + String method, + Exception e, [ + StackTrace? stackTrace, + ]) { + if (e is InvalidRequest) { + logger.w("gRPC $method rejected: ${e.message}"); + } else { + logger.e( + "gRPC $method failed", + error: e, + stackTrace: stackTrace, + ); + } + return grpc.GrpcError.unknown(e.toString()); + } - Future _handleExceptions(Future Function() f) async { + Future _handleExceptions( + String method, + Future Function() f, + ) async { + logger.d("gRPC $method received"); try { - return await f(); - } on Exception catch (e) { - throw _wrapException(e); + final result = await f(); + logger.d("gRPC $method completed"); + return result; + } on Exception catch (e, stackTrace) { + throw _wrapException(method, e, stackTrace); } } - Future _handleEmpty(Future Function() f) async { - await _handleExceptions(f); + Future _handleEmpty( + String method, + Future Function() f, + ) async { + await _handleExceptions(method, f); return pb.Empty(); } @@ -43,11 +66,11 @@ class FrostNoosphereService extends pb.NoosphereServiceBase { grpc.ServiceCall call, pb.LoginRequest request, ) => - _handleExceptions(() async { + _handleExceptions("login", () async { final resp = await api.login( - groupFingerprint: _bytes(request.groupFingerprint), + groupFingerprint: common.bytes(request.groupFingerprint), participantId: Identifier.fromBytes( - _bytes(request.participantId), + common.bytes(request.participantId), ), protocolVersion: request.protocolVersion, ); @@ -60,11 +83,11 @@ class FrostNoosphereService extends pb.NoosphereServiceBase { grpc.ServiceCall call, pb.SignedAuthChallenge request, ) => - _handleExceptions(() async { + _handleExceptions("respondToChallenge", () async { final resp = await api.respondToChallenge( Signed( - obj: AuthChallenge.fromBytes(_bytes(request.challenge)), - signature: cl.SchnorrSignature(_bytes(request.signature)), + obj: AuthChallenge.fromBytes(common.bytes(request.challenge)), + signature: cl.SchnorrSignature(common.bytes(request.signature)), ), ); @@ -76,24 +99,40 @@ class FrostNoosphereService extends pb.NoosphereServiceBase { grpc.ServiceCall call, pb.Bytes request, ) { - final sessionId = _sid(request.data); + final sessionId = common.sid(request.data); late final ClientSession session; try { session = api.getSession(sessionId); - } on Exception catch (e) { - throw _wrapException(e); + } on Exception catch (e, stackTrace) { + throw _wrapException("fetchEventStream", e, stackTrace); } + logger.d( + "gRPC fetchEventStream opened for participant ${session.participantId}", + ); + // sendTrailers is not always called automatically when the stream ends // despite the documentation. // Without calling this, the grpc stream may hang and never close. final controller = StreamController( - onCancel: () => call.sendTrailers(), + onCancel: () { + logger.d( + "gRPC fetchEventStream canceled for participant " + "${session.participantId}", + ); + call.sendTrailers(); + }, ); // When upstream stream is done, cancel this one controller.addStream(session.eventController.stream).then( - (_) => controller.close(), + (_) { + logger.d( + "gRPC fetchEventStream closed for participant " + "${session.participantId}", ); + return controller.close(); + }, + ); // Pass across all events return controller.stream.map( @@ -124,8 +163,8 @@ class FrostNoosphereService extends pb.NoosphereServiceBase { grpc.ServiceCall call, pb.Bytes request, ) => - _handleExceptions(() async { - final resp = await api.extendSession(_sid(request.data)); + _handleExceptions("extendSession", () async { + final resp = await api.extendSession(common.sid(request.data)); return _returnWritable(resp); }); @@ -135,14 +174,15 @@ class FrostNoosphereService extends pb.NoosphereServiceBase { pb.DkgRequest request, ) => _handleEmpty( + "requestNewDkg", () => api.requestNewDkg( - sid: _sid(request.sid), + sid: common.sid(request.sid), signedDetails: Signed.fromBytes( - _bytes(request.signedDetails), + common.bytes(request.signedDetails), (reader) => NewDkgDetails.fromReader(reader), ), commitment: DkgPublicCommitment.fromBytes( - _bytes(request.commitment), + common.bytes(request.commitment), ), ), ); @@ -153,7 +193,8 @@ class FrostNoosphereService extends pb.NoosphereServiceBase { pb.DkgToReject request, ) => _handleEmpty( - () => api.rejectDkg(sid: _sid(request.sid), name: request.name), + "rejectDkg", + () => api.rejectDkg(sid: common.sid(request.sid), name: request.name), ); @override @@ -162,11 +203,12 @@ class FrostNoosphereService extends pb.NoosphereServiceBase { pb.DkgCommitment request, ) => _handleEmpty( + "submitDkgCommitment", () => api.submitDkgCommitment( - sid: _sid(request.sid), + sid: common.sid(request.sid), name: request.name, commitment: DkgPublicCommitment.fromBytes( - _bytes(request.commitment), + common.bytes(request.commitment), ), ), ); @@ -177,16 +219,17 @@ class FrostNoosphereService extends pb.NoosphereServiceBase { pb.DkgRound2 request, ) => _handleEmpty( + "submitDkgRound2", () => api.submitDkgRound2( - sid: _sid(request.sid), + sid: common.sid(request.sid), name: request.name, commitmentSetSignature: cl.SchnorrSignature( - _bytes(request.commitmentSetSignature), + common.bytes(request.commitmentSetSignature), ), secrets: { for (final secret in request.secrets) - Identifier.fromBytes(_bytes(secret.id)): DkgEncryptedSecret( - ECCiphertext.fromBytes(_bytes(secret.secret)), + Identifier.fromBytes(common.bytes(secret.id)): DkgEncryptedSecret( + ECCiphertext.fromBytes(common.bytes(secret.secret)), ), }, ), @@ -198,11 +241,12 @@ class FrostNoosphereService extends pb.NoosphereServiceBase { pb.DkgAcks request, ) => _handleEmpty( + "sendDkgAcks", () => api.sendDkgAcks( - sid: _sid(request.sid), + sid: common.sid(request.sid), acks: request.acks .map( - (ack) => SignedDkgAck.fromBytes(_bytes(ack)), + (ack) => SignedDkgAck.fromBytes(common.bytes(ack)), ) .toSet(), ), @@ -213,12 +257,12 @@ class FrostNoosphereService extends pb.NoosphereServiceBase { grpc.ServiceCall call, pb.DkgAckRequest request, ) => - _handleExceptions(() async { + _handleExceptions("requestDkgAcks", () async { final resp = await api.requestDkgAcks( - sid: _sid(request.sid), + sid: common.sid(request.sid), requests: request.requests .map( - (request) => DkgAckRequest.fromBytes(_bytes(request)), + (request) => DkgAckRequest.fromBytes(common.bytes(request)), ) .toSet(), ); @@ -232,20 +276,21 @@ class FrostNoosphereService extends pb.NoosphereServiceBase { pb.SignaturesRequest request, ) => _handleEmpty( + "requestSignatures", () => api.requestSignatures( - sid: _sid(request.sid), + sid: common.sid(request.sid), keys: request.keys .map( - (key) => AggregateKeyInfo.fromBytes(_bytes(key)), + (key) => AggregateKeyInfo.fromBytes(common.bytes(key)), ) .toSet(), signedDetails: Signed.fromBytes( - _bytes(request.signedDetails), + common.bytes(request.signedDetails), (reader) => SignaturesRequestDetails.fromReader(reader), ), commitments: request.commitments .map( - (commitment) => SigningCommitment.fromBytes(_bytes(commitment)), + (commitment) => SigningCommitment.fromBytes(common.bytes(commitment)), ) .toList(), ), @@ -257,9 +302,10 @@ class FrostNoosphereService extends pb.NoosphereServiceBase { pb.SignaturesRejection request, ) => _handleEmpty( + "rejectSignaturesRequest", () => api.rejectSignaturesRequest( - sid: _sid(request.sid), - reqId: _sigReqId(request.reqId), + sid: common.sid(request.sid), + reqId: common.sigReqId(request.reqId), ), ); @@ -268,13 +314,13 @@ class FrostNoosphereService extends pb.NoosphereServiceBase { grpc.ServiceCall call, pb.SignaturesReplies request, ) => - _handleExceptions(() async { + _handleExceptions("submitSignatureReplies", () async { final resp = await api.submitSignatureReplies( - sid: _sid(request.sid), - reqId: _sigReqId(request.reqId), + sid: common.sid(request.sid), + reqId: common.sigReqId(request.reqId), replies: request.replies .map( - (reply) => SignatureReply.fromBytes(_bytes(reply)), + (reply) => SignatureReply.fromBytes(common.bytes(reply)), ) .toList(), ); @@ -296,14 +342,14 @@ class FrostNoosphereService extends pb.NoosphereServiceBase { grpc.ServiceCall call, pb.SecretShare request, ) => - _handleExceptions(() async { + _handleExceptions("shareSecretShare", () async { final resp = await api.shareSecretShare( - sid: _sid(request.sid), - groupKey: cl.ECCompressedPublicKey(_bytes(request.groupKey)), + sid: common.sid(request.sid), + groupKey: cl.ECCompressedPublicKey(common.bytes(request.groupKey)), encryptedSecrets: { for (final secret in request.secrets) - Identifier.fromBytes(_bytes(secret.id)): EncryptedKeyShare( - ECCiphertext.fromBytes(_bytes(secret.share)), + Identifier.fromBytes(common.bytes(secret.id)): EncryptedKeyShare( + ECCiphertext.fromBytes(common.bytes(secret.share)), ), }, ); @@ -317,10 +363,11 @@ class FrostNoosphereService extends pb.NoosphereServiceBase { pb.ConstructedKey request, ) => _handleEmpty( + "ackKeyConstructed", () => api.ackKeyConstructed( - sid: _sid(request.sid), + sid: common.sid(request.sid), constructedKey: Signed.fromBytes( - _bytes(request.constructedKey), + common.bytes(request.constructedKey), (reader) => KeyWasConstructed.fromReader(reader), ), ), diff --git a/lib/src/logging.dart b/lib/src/logging.dart new file mode 100644 index 0000000..e7e92ba --- /dev/null +++ b/lib/src/logging.dart @@ -0,0 +1,19 @@ +export 'package:logger/logger.dart' + show Level, Logger, LogFilter, LogOutput, LogPrinter; + +import 'package:logger/logger.dart'; + +/// Creates the default logger used by server components when no logger is +/// provided. +Logger createNoosphereRoastServerLogger({ + Level level = Level.info, + LogFilter? filter, + LogPrinter? printer, + LogOutput? output, +}) => + Logger( + filter: filter ?? ProductionFilter(), + printer: printer ?? SimplePrinter(printTime: true, colors: false), + output: output, + level: level, + ); diff --git a/lib/src/noosphere_roast_server_base.dart b/lib/src/noosphere_roast_server_base.dart index 818c6f7..ac54ba0 100644 --- a/lib/src/noosphere_roast_server_base.dart +++ b/lib/src/noosphere_roast_server_base.dart @@ -3,4 +3,7 @@ export "package:noosphere_roast_client/noosphere_roast_client.dart"; export "config/grpc.dart"; export "config/server.dart"; export "grpc.dart"; +export "logging.dart"; +export "rest.dart"; export "server/api_handler.dart"; +export "server/synchronized_api_handler.dart"; diff --git a/lib/src/rest.dart b/lib/src/rest.dart new file mode 100644 index 0000000..2dd34da --- /dev/null +++ b/lib/src/rest.dart @@ -0,0 +1,495 @@ +import 'dart:async'; +import 'dart:convert'; +import 'dart:io'; +import 'dart:typed_data'; +import 'package:coinlib/coinlib.dart' as cl; +import 'package:noosphere_roast_client/noosphere_roast_client.dart'; +import 'package:noosphere_roast_server/src/common.dart' as common; +import 'package:noosphere_roast_server/src/logging.dart'; +import 'package:noosphere_roast_server/src/server/api_handler.dart'; +import 'package:shelf/shelf.dart'; +import 'package:shelf/shelf_io.dart' as shelf_io; +import 'package:shelf_router/shelf_router.dart'; +import 'package:shelf_web_socket/shelf_web_socket.dart'; +import 'package:web_socket_channel/web_socket_channel.dart'; + +String _base64Encode(List bytes) => base64Encode(bytes); +String _base64UrlEncode(List bytes) => + base64UrlEncode(bytes).replaceAll('=', ''); + +Uint8List _decodeBytes(String value) { + final base64Value = value.replaceAll('-', '+').replaceAll('_', '/'); + final padding = (4 - base64Value.length % 4) % 4; + return base64Decode(base64Value.padRight(base64Value.length + padding, '=')); +} + +Map _corsHeaders(String allowOrigin) => { + 'access-control-allow-origin': allowOrigin, + 'access-control-allow-methods': 'GET, POST, OPTIONS', + 'access-control-allow-headers': 'content-type', + 'access-control-max-age': '86400', + }; + +Middleware restWebSocketCors({ + String allowOrigin = '*', +}) => + (innerHandler) => (request) async { + final headers = _corsHeaders(allowOrigin); + if (request.method == 'OPTIONS') { + return Response.ok('', headers: headers); + } + + final response = await innerHandler(request); + return response.change(headers: {...response.headers, ...headers}); + }; + +class RestWebSocketNoosphereService { + final ServerApiHandler api; + final String? allowOrigin; + final Logger logger; + + RestWebSocketNoosphereService({ + required this.api, + this.allowOrigin = '*', + Logger? logger, + }) : logger = logger ?? api.logger; + + Handler get handler { + final router = Router() + ..post('/login', _login) + ..post('/respond-to-challenge', _respondToChallenge) + ..post('/extend-session', _extendSession) + ..post('/dkg/new', _requestNewDkg) + ..post('/dkg/reject', _rejectDkg) + ..post('/dkg/commitment', _submitDkgCommitment) + ..post('/dkg/round2', _submitDkgRound2) + ..post('/dkg/acks', _sendDkgAcks) + ..post('/dkg/request-acks', _requestDkgAcks) + ..post('/signatures/request', _requestSignatures) + ..post('/signatures/reject', _rejectSignaturesRequest) + ..post('/signatures/replies', _submitSignatureReplies) + ..post('/secret-share', _shareSecretShare) + ..post('/key-constructed/ack', _ackKeyConstructed) + ..get('/sessions//events', _fetchEventWebSocket); + + var pipeline = const Pipeline(); + final allowOrigin = this.allowOrigin; + if (allowOrigin != null) { + pipeline = pipeline.addMiddleware( + restWebSocketCors(allowOrigin: allowOrigin), + ); + } + return pipeline.addHandler(router.call); + } + + Future serve({ + Object address = 'localhost', + int port = 8080, + }) => + shelf_io.serve(handler, address, port); + + Future _login(Request request) => + _handleRequest(request, logger, () async { + final json = await _readJson(request); + final resp = await api.login( + groupFingerprint: _fieldBytes(json, 'groupFingerprint'), + participantId: + Identifier.fromBytes(_fieldBytes(json, 'participantId')), + protocolVersion: _optionalInt(json, 'protocolVersion') ?? + ServerApiHandler.currentProtocolVersion, + ); + return _bytesResponse(resp.toBytes()); + }); + + Future _respondToChallenge(Request request) => _handleRequest( + request, + logger, + () async { + final json = await _readJson(request); + final resp = await api.respondToChallenge( + Signed( + obj: AuthChallenge.fromBytes(_fieldBytes(json, 'challenge')), + signature: cl.SchnorrSignature(_fieldBytes(json, 'signature')), + ), + ); + return _bytesResponse(resp.toBytes()); + }, + ); + + Future _extendSession(Request request) => + _handleRequest(request, logger, () async { + final json = await _readJson(request); + final resp = await api.extendSession(common.sid(_fieldBytes(json, 'sid'))); + return _bytesResponse(resp.toBytes()); + }); + + Future _requestNewDkg(Request request) => + _handleEmpty(request, logger, () async { + final json = await _readJson(request); + await api.requestNewDkg( + sid: common.sid(_fieldBytes(json, 'sid')), + signedDetails: Signed.fromBytes( + _fieldBytes(json, 'signedDetails'), + (reader) => NewDkgDetails.fromReader(reader), + ), + commitment: DkgPublicCommitment.fromBytes( + _fieldBytes(json, 'commitment'), + ), + ); + }); + + Future _rejectDkg(Request request) => + _handleEmpty(request, logger, () async { + final json = await _readJson(request); + await api.rejectDkg( + sid: common.sid(_fieldBytes(json, 'sid')), + name: _fieldString(json, 'name'), + ); + }); + + Future _submitDkgCommitment(Request request) => + _handleEmpty(request, logger, () async { + final json = await _readJson(request); + await api.submitDkgCommitment( + sid: common.sid(_fieldBytes(json, 'sid')), + name: _fieldString(json, 'name'), + commitment: DkgPublicCommitment.fromBytes( + _fieldBytes(json, 'commitment'), + ), + ); + }); + + Future _submitDkgRound2(Request request) => + _handleEmpty(request, logger, () async { + final json = await _readJson(request); + await api.submitDkgRound2( + sid: common.sid(_fieldBytes(json, 'sid')), + name: _fieldString(json, 'name'), + commitmentSetSignature: cl.SchnorrSignature( + _fieldBytes(json, 'commitmentSetSignature'), + ), + secrets: { + for (final secret in _fieldList(json, 'secrets')) + Identifier.fromBytes(_fieldBytes(secret, 'id')): + DkgEncryptedSecret( + ECCiphertext.fromBytes(_fieldBytes(secret, 'secret')), + ), + }, + ); + }); + + Future _sendDkgAcks(Request request) => + _handleEmpty(request, logger, () async { + final json = await _readJson(request); + await api.sendDkgAcks( + sid: common.sid(_fieldBytes(json, 'sid')), + acks: _fieldBytesList( + json, + 'acks', + ).map(SignedDkgAck.fromBytes).toSet(), + ); + }); + + Future _requestDkgAcks(Request request) => + _handleRequest(request, logger, () async { + final json = await _readJson(request); + final resp = await api.requestDkgAcks( + sid: common.sid(_fieldBytes(json, 'sid')), + requests: _fieldBytesList( + json, + 'requests', + ).map(DkgAckRequest.fromBytes).toSet(), + ); + return _repeatedBytesResponse(resp.map((ack) => ack.toBytes())); + }); + + Future _requestSignatures(Request request) => + _handleEmpty(request, logger, () async { + final json = await _readJson(request); + await api.requestSignatures( + sid: common.sid(_fieldBytes(json, 'sid')), + keys: _fieldBytesList( + json, + 'keys', + ).map(AggregateKeyInfo.fromBytes).toSet(), + signedDetails: Signed.fromBytes( + _fieldBytes(json, 'signedDetails'), + (reader) => SignaturesRequestDetails.fromReader(reader), + ), + commitments: _fieldBytesList( + json, + 'commitments', + ).map(SigningCommitment.fromBytes).toList(), + ); + }); + + Future _rejectSignaturesRequest(Request request) => + _handleEmpty(request, logger, () async { + final json = await _readJson(request); + await api.rejectSignaturesRequest( + sid: common.sid(_fieldBytes(json, 'sid')), + reqId: common.sigReqId(_fieldBytes(json, 'reqId')), + ); + }); + + Future _submitSignatureReplies(Request request) => + _handleRequest(request, logger, () async { + final json = await _readJson(request); + final resp = await api.submitSignatureReplies( + sid: common.sid(_fieldBytes(json, 'sid')), + reqId: common.sigReqId(_fieldBytes(json, 'reqId')), + replies: _fieldBytesList( + json, + 'replies', + ).map(SignatureReply.fromBytes).toList(), + ); + + return _jsonResponse({ + 'type': switch (resp) { + SignatureNewRoundsResponse() => 'new_round', + SignaturesCompleteResponse() => 'complete', + null => 'empty', + }, + 'data': resp == null ? null : _base64Encode(resp.toBytes()), + }); + }); + + Future _shareSecretShare(Request request) => + _handleRequest(request, logger, () async { + final json = await _readJson(request); + final resp = await api.shareSecretShare( + sid: common.sid(_fieldBytes(json, 'sid')), + groupKey: cl.ECCompressedPublicKey(_fieldBytes(json, 'groupKey')), + encryptedSecrets: { + for (final secret in _fieldList(json, 'secrets')) + Identifier.fromBytes(_fieldBytes(secret, 'id')): + EncryptedKeyShare( + ECCiphertext.fromBytes(_fieldBytes(secret, 'share')), + ), + }, + ); + return _repeatedBytesResponse(resp.map((ev) => ev.toBytes())); + }); + + Future _ackKeyConstructed(Request request) => + _handleEmpty(request, logger, () async { + final json = await _readJson(request); + await api.ackKeyConstructed( + sid: common.sid(_fieldBytes(json, 'sid')), + constructedKey: Signed.fromBytes( + _fieldBytes(json, 'constructedKey'), + (reader) => KeyWasConstructed.fromReader(reader), + ), + ); + }); + + FutureOr _fetchEventWebSocket(Request request, String sid) { + final description = _requestDescription(request); + logger.d("REST $description received"); + try { + final session = api.getSession(common.sid(_decodeBytes(sid))); + final handler = webSocketHandler( + (WebSocketChannel webSocket, String? _) { + logger.d("REST $description opened"); + final eventSubscription = session.eventController.stream.listen( + (event) => webSocket.sink.add(_webSocketEvent(event, logger)), + onDone: () { + unawaited(webSocket.sink.close(WebSocketStatus.normalClosure)); + }, + onError: (Object e, StackTrace stackTrace) { + logger.e( + "REST $description event stream failed", + error: e, + stackTrace: stackTrace, + ); + unawaited( + webSocket.sink.close(WebSocketStatus.internalServerError), + ); + }, + cancelOnError: true, + ); + webSocket.stream.listen( + (_) {}, + onDone: () { + unawaited(eventSubscription.cancel()); + logger.d("REST $description closed"); + }, + onError: (Object e) { + logger.w( + "REST $description socket failed: $e", + ); + unawaited(eventSubscription.cancel()); + }, + cancelOnError: true, + ); + }, + allowedOrigins: _webSocketAllowedOrigins, + ); + return handler(request); + } on InvalidRequest catch (e) { + logger.w( + "REST $description rejected: ${e.message}", + ); + return _jsonResponse({'error': e.message}, status: 400); + } on FormatException catch (e) { + logger.w( + "REST $description rejected: ${e.message}", + ); + return _jsonResponse({'error': e.message}, status: 400); + } on HijackException { + rethrow; + } on Exception catch (e, stackTrace) { + logger.e( + "REST $description failed", + error: e, + stackTrace: stackTrace, + ); + return _jsonResponse({'error': 'Internal server error'}, status: 500); + } + } + + Iterable? get _webSocketAllowedOrigins { + final allowOrigin = this.allowOrigin; + if (allowOrigin == null || allowOrigin == '*') return null; + return [allowOrigin]; + } +} + +Response _rejectResponse(String description, String message, Logger logger) { + logger.w("REST $description rejected: $message"); + return _jsonResponse({'error': message}, status: 400); +} + +Future _handleEmpty( + Request request, + Logger logger, + Future Function() action, +) async { + return _handleRequest(request, logger, () async { + await action(); + return _jsonResponse({}); + }); +} + +Future _handleRequest( + Request request, + Logger logger, + Future Function() action, +) async { + final description = _requestDescription(request); + logger.d("REST $description received"); + try { + final response = await action(); + logger.d("REST $description completed"); + return response; + } on InvalidRequest catch (e) { + return _rejectResponse(description, e.message, logger); + } on FormatException catch (e) { + return _rejectResponse(description, e.message, logger); + } on Exception catch (e, stackTrace) { + logger.e( + "REST $description failed", + error: e, + stackTrace: stackTrace, + ); + return _jsonResponse({'error': 'Internal server error'}, status: 500); + } +} + +String _requestDescription(Request request) { + final path = request.url.path.isEmpty ? '/' : '/${request.url.path}'; + return '${request.method} $path'; +} + +Future> _readJson(Request request) async { + final body = await request.readAsString(); + final decoded = jsonDecode(body); + if (decoded is! Map) { + throw const FormatException('Expected JSON object'); + } + return decoded; +} + +Response _jsonResponse(Object value, {int status = 200}) => Response( + status, + body: jsonEncode(value), + headers: {'content-type': 'application/json'}, + ); + +Response _bytesResponse(List bytes) => + _jsonResponse({'data': _base64Encode(bytes)}); + +Response _repeatedBytesResponse(Iterable> bytes) => + _jsonResponse({'data': bytes.map(_base64Encode).toList()}); + +String _fieldString(Map json, String name) { + final value = json[name]; + if (value is! String) { + throw FormatException('Expected "$name" to be a string'); + } + return value; +} + +int? _optionalInt(Map json, String name) { + final value = json[name]; + if (value == null) return null; + if (value is! int) throw FormatException('Expected "$name" to be an int'); + return value; +} + +Uint8List _fieldBytes(Map json, String name) => + _decodeBytes(_fieldString(json, name)); + +List> _fieldList(Map json, String name) { + final value = json[name]; + if (value is! List) throw FormatException('Expected "$name" to be a list'); + return value.map((entry) { + if (entry is! Map) { + throw FormatException('Expected "$name" entries to be objects'); + } + return { + for (final mapEntry in entry.entries) + if (mapEntry.key is String) mapEntry.key as String: mapEntry.value, + }; + }).toList(); +} + +Iterable _fieldBytesList(Map json, String name) { + final value = json[name]; + if (value is! List) throw FormatException('Expected "$name" to be a list'); + return value.map((entry) { + if (entry is! String) { + throw FormatException('Expected "$name" entries to be strings'); + } + return _decodeBytes(entry); + }); +} + +String _webSocketEvent(Event event, Logger logger) { + final type = _eventType(event); + logger.d("REST WebSocket sent $type"); + return jsonEncode({ + 'type': type, + 'data': _base64Encode(event.toBytes()), + }); +} + +String _eventType(Event event) => switch (event) { + ParticipantStatusEvent() => 'participant_status', + NewDkgEvent() => 'new_dkg', + DkgCommitmentEvent() => 'dkg_commitment', + DkgRejectEvent() => 'dkg_reject', + DkgRound2ShareEvent() => 'dkg_round2_share', + DkgAckEvent() => 'dkg_ack', + DkgAckRequestEvent() => 'dkg_ack_request', + SignaturesRequestEvent() => 'signatures_request', + SignatureNewRoundsEvent() => 'signature_new_rounds', + SignaturesCompleteEvent() => 'signatures_complete', + SignaturesFailureEvent() => 'signatures_failure', + SecretShareEvent() => 'secret_share', + ConstructedKeyEvent() => 'constructed_key', + KeepaliveEvent() => 'keepalive', + }; + +String restWebSocketSessionPath(SessionID sid) => + '/sessions/${_base64UrlEncode(sid.n)}/events'; diff --git a/lib/src/server/api_handler.dart b/lib/src/server/api_handler.dart index 874be77..fa1777a 100644 --- a/lib/src/server/api_handler.dart +++ b/lib/src/server/api_handler.dart @@ -5,6 +5,7 @@ import 'package:collection/collection.dart'; import 'package:coinlib/coinlib.dart' as cl; import 'package:noosphere_roast_client/noosphere_roast_client.dart'; import 'package:noosphere_roast_server/src/config/server.dart'; +import 'package:noosphere_roast_server/src/logging.dart'; import 'package:noosphere_roast_server/src/server/state/key_sharing.dart'; import 'state/signatures_coordination.dart'; import 'state/client_session.dart'; @@ -21,15 +22,21 @@ class ServerApiHandler implements ApiRequestInterface { static const currentProtocolVersion = 2; final ServerConfig config; - final ServerState state; + late final ServerState state; + late final Logger logger; final DateTime startTime = DateTime.now(); - /// Creates a backend API handler with the [config]. A blank [state] will be - /// created if not provided. + /// Creates a backend API handler with the [config]. + /// + /// A blank [state] and default [logger] will be created if not provided. ServerApiHandler({ required this.config, ServerState? state, - }) : state = state ?? ServerState(); + Logger? logger, + }) { + this.logger = logger ?? createNoosphereRoastServerLogger(); + this.state = state ?? ServerState(logger: this.logger); + } int get _participantN => config.group.participants.length; @@ -88,6 +95,10 @@ class ServerApiHandler implements ApiRequestInterface { expiry: expiry, ); + logger.i( + "Issued auth challenge for participant $participantId", + ); + return ExpirableAuthChallengeResponse(challenge: challenge, expiry: expiry); } @@ -157,6 +168,8 @@ class ServerApiHandler implements ApiRequestInterface { }); } + logger.i("Participant logged in: $pid"); + return LoginCompleteResponse( id: sessionId, expiry: expiry, @@ -271,6 +284,11 @@ class ServerApiHandler implements ApiRequestInterface { commitments: commitments, ); + logger.i( + "DKG requested: name=${details.name} creator=${session.participantId} " + "threshold=${details.threshold}", + ); + // Broadcast to other participants state.sendEventToOthers(dkgEvent, sid); } @@ -279,6 +297,10 @@ class ServerApiHandler implements ApiRequestInterface { Future rejectDkg({required SessionID sid, required String name}) async { final participantId = getSession(sid).participantId; if (state.nameToDkg.remove(name) != null) { + logger.i( + "DKG rejected: name=$name participant=$participantId", + ); + // Send an event to all other participants that the DKG was removed state.sendEventToOthers( DkgRejectEvent(name: name, participant: participantId), @@ -313,6 +335,9 @@ class ServerApiHandler implements ApiRequestInterface { dkg.round = DkgRound2State( expectedHash: dkg.details.obj.hashWithCommitments(commitmentSet), ); + logger.i( + "DKG advanced to round 2: name=$name commitments=${commitments.length}", + ); } // Send commitment to other participants @@ -372,6 +397,7 @@ class ServerApiHandler implements ApiRequestInterface { if (round.participantsProvided.length == _participantN - 1) { // Remove DKG state.nameToDkg.remove(name); + logger.i("DKG completed: name=$name"); // No details of the key are stored on the server as only the participants // can generate the public information at this point. } else { @@ -417,6 +443,8 @@ class ServerApiHandler implements ApiRequestInterface { // Do not send events if there are no new ACKs if (newAcks.isEmpty) return; + logger.i("DKG acknowledgements received: ${newAcks.length}"); + // Send ACKs to participants, ensuring that their own ACKs aren't sent // Do not send to calling participant for (final session in state.clientSessions.values.where( @@ -485,6 +513,10 @@ class ServerApiHandler implements ApiRequestInterface { } if (need.isNotEmpty) { + logger.d( + "Requested missing DKG acknowledgements: ${need.length}", + ); + // Send DkgAckRequestEvents for missing ACKs state.sendEventToOthers(DkgAckRequestEvent(need), sid); } @@ -557,6 +589,11 @@ class ServerApiHandler implements ApiRequestInterface { ), sid, ); + + logger.i( + "Signatures requested: id=${details.id.toHex()} creator=$pid " + "signatures=$numSigs", + ); } void _checkSigReqFail(SignaturesCoordinationState sigReqState) { @@ -571,6 +608,10 @@ class ServerApiHandler implements ApiRequestInterface { if (available < maxThreshold) { // Cannot sign one of the signatures as threshold is too high final id = sigReqState.details.obj.id; + logger.w( + "Signatures request failed: id=${id.toHex()} available=$available " + "required=$maxThreshold", + ); state.sendEventToAll(SignaturesFailureEvent(id)); state.sigRequests.remove(id); } @@ -592,6 +633,9 @@ class ServerApiHandler implements ApiRequestInterface { if (sigReq.malicious.contains(pid)) return; sigReq.rejectors.add(pid); + logger.i( + "Signatures request rejected: id=${reqId.toHex()} participant=$pid", + ); _checkSigReqFail(sigReq); } @@ -611,6 +655,10 @@ class ServerApiHandler implements ApiRequestInterface { void throwMalicious(InvalidRequest exp) { sigReq.malicious.add(pid); + logger.w( + "Participant marked malicious for signatures request: " + "id=${reqId.toHex()} participant=$pid reason=${exp.message}", + ); _checkSigReqFail(sigReq); throw exp; } @@ -768,12 +816,22 @@ class ServerApiHandler implements ApiRequestInterface { sid, ); + logger.i( + "Signatures request completed: id=${reqId.toHex()} " + "signatures=${signatures.length}", + ); + return SignaturesCompleteResponse(signatures); } // If there are any new rounds, return them and send events to round // participants if (newRounds.isNotEmpty) { + logger.d( + "Signature rounds started: id=${reqId.toHex()} " + "participants=${newRounds.length}", + ); + for (final id in newRounds.keys.where((id) => id != pid)) { state.participantToSession[id]?.sendEvent( SignatureNewRoundsEvent(reqId: reqId, rounds: newRounds[id]!), @@ -816,15 +874,22 @@ class ServerApiHandler implements ApiRequestInterface { // events. final secrets = state.secretSharesForKey(groupKey); + var addedShares = 0; for (final MapEntry(key: id, value: share) in encryptedSecrets.entries) { if (secrets.maybeAddShare(pid, id, share)) { + addedShares++; state.participantToSession[id]?.sendEvent( SecretShareEvent(sender: pid, keyShare: share, groupKey: groupKey), ); } } + logger.i( + "Secret shares received: sender=$pid receivers=${encryptedSecrets.length} " + "new=$addedShares", + ); + // Return cached ConstructedKeyEvents for unneeded secrets return secrets.eventsForCompleted(encryptedSecrets.keys); } @@ -859,12 +924,21 @@ class ServerApiHandler implements ApiRequestInterface { // Send event to other participants state.sendEventToOthers(event, sid); + + logger.i( + "Constructed key acknowledged: participant=$pid", + ); } /// Closes all client session streams - Future shutdown() => Future.wait( - state.clientSessions.values.map( - (session) => session.eventController.close(), - ), - ); + Future shutdown() { + logger.i( + "Shutting down API handler: sessions=${state.clientSessions.values.length}", + ); + return Future.wait( + state.clientSessions.values.map( + (session) => session.eventController.close(), + ), + ); + } } diff --git a/lib/src/server/state/state.dart b/lib/src/server/state/state.dart index e6524a5..8996629 100644 --- a/lib/src/server/state/state.dart +++ b/lib/src/server/state/state.dart @@ -1,6 +1,7 @@ import 'package:coinlib/coinlib.dart' as cl; import 'package:noosphere_roast_client/common.dart'; import 'package:noosphere_roast_client/noosphere_roast_client.dart'; +import 'package:noosphere_roast_server/src/logging.dart'; import 'client_session.dart'; import 'dkg.dart'; import 'key_sharing.dart'; @@ -45,6 +46,7 @@ class CompletedSignatures implements Expirable { } class ServerState { + final Logger logger; final challenges = ExpirableMap(); late final ExpirableMap clientSessions; final participantToSession = ExpirableMap(); @@ -59,13 +61,19 @@ class ServerState { /// participants final Map secretShares = {}; - ServerState() { + ServerState({ + required this.logger, + }) { clientSessions = ExpirableMap( onExpired: (_, session) => onEndSession(session), ); } void onEndSession(ClientSession session) { + logger.i( + "Participant session ended: ${session.participantId}", + ); + // Reset DKGs to round 1 as all participants need to remain online to // complete them for (final dkg in nameToDkg.values) { @@ -91,8 +99,15 @@ class ServerState { ); void sendEventToAll(Event e, {List exclude = const []}) { - for (final session in clientSessions.values) { - if (!exclude.contains(session.sessionID)) session.sendEvent(e); + final recipients = clientSessions.values + .where((session) => !exclude.contains(session.sessionID)) + .toList(); + logger.d( + "Broadcasting ${e.runtimeType} to ${recipients.length}/" + "${clientSessions.values.length} sessions", + ); + for (final session in recipients) { + session.sendEvent(e); } } diff --git a/lib/src/server/synchronized_api_handler.dart b/lib/src/server/synchronized_api_handler.dart new file mode 100644 index 0000000..b5a9425 --- /dev/null +++ b/lib/src/server/synchronized_api_handler.dart @@ -0,0 +1,195 @@ +import 'dart:async'; +import 'dart:typed_data'; +import 'package:coinlib/coinlib.dart' as cl; +import 'package:noosphere_roast_client/noosphere_roast_client.dart'; +import 'package:noosphere_roast_server/src/server/api_handler.dart'; + +class _ApiCallQueue { + Future _tail = Future.value(); + + Future run(Future Function() action) { + final previous = _tail; + final completer = Completer(); + + _tail = previous.catchError((_) {}).then((_) async { + try { + completer.complete(await action()); + } catch (e, st) { + completer.completeError(e, st); + } + }); + + return completer.future; + } +} + +/// A [ServerApiHandler] that serializes state-mutating API calls. +/// +/// Use one shared instance of this class when exposing the same coordinator +/// through multiple transports, such as gRPC for desktop clients and +/// REST/WebSocket for web clients. +class SynchronizedServerApiHandler extends ServerApiHandler { + final _queue = _ApiCallQueue(); + + SynchronizedServerApiHandler({ + required super.config, + super.state, + super.logger, + }); + + @override + Future login({ + required Uint8List groupFingerprint, + required Identifier participantId, + int protocolVersion = ServerApiHandler.currentProtocolVersion, + }) => + _queue.run( + () => super.login( + groupFingerprint: groupFingerprint, + participantId: participantId, + protocolVersion: protocolVersion, + ), + ); + + @override + Future respondToChallenge( + Signed signedChallenge, + ) => + _queue.run(() => super.respondToChallenge(signedChallenge)); + + @override + Future extendSession(SessionID sid) => + _queue.run(() => super.extendSession(sid)); + + @override + Future requestNewDkg({ + required SessionID sid, + required Signed signedDetails, + required DkgPublicCommitment commitment, + }) => + _queue.run( + () => super.requestNewDkg( + sid: sid, + signedDetails: signedDetails, + commitment: commitment, + ), + ); + + @override + Future rejectDkg({ + required SessionID sid, + required String name, + }) => + _queue.run(() => super.rejectDkg(sid: sid, name: name)); + + @override + Future submitDkgCommitment({ + required SessionID sid, + required String name, + required DkgPublicCommitment commitment, + }) => + _queue.run( + () => super.submitDkgCommitment( + sid: sid, + name: name, + commitment: commitment, + ), + ); + + @override + Future submitDkgRound2({ + required SessionID sid, + required String name, + required cl.SchnorrSignature commitmentSetSignature, + required Map secrets, + }) => + _queue.run( + () => super.submitDkgRound2( + sid: sid, + name: name, + commitmentSetSignature: commitmentSetSignature, + secrets: secrets, + ), + ); + + @override + Future sendDkgAcks({ + required SessionID sid, + required Set acks, + }) => + _queue.run(() => super.sendDkgAcks(sid: sid, acks: acks)); + + @override + Future> requestDkgAcks({ + required SessionID sid, + required Set requests, + }) => + _queue.run( + () => super.requestDkgAcks(sid: sid, requests: requests), + ); + + @override + Future requestSignatures({ + required SessionID sid, + required Set keys, + required Signed signedDetails, + required List commitments, + }) => + _queue.run( + () => super.requestSignatures( + sid: sid, + keys: keys, + signedDetails: signedDetails, + commitments: commitments, + ), + ); + + @override + Future rejectSignaturesRequest({ + required SessionID sid, + required SignaturesRequestId reqId, + }) => + _queue.run( + () => super.rejectSignaturesRequest(sid: sid, reqId: reqId), + ); + + @override + Future submitSignatureReplies({ + required SessionID sid, + required SignaturesRequestId reqId, + required List replies, + }) => + _queue.run( + () => super.submitSignatureReplies( + sid: sid, + reqId: reqId, + replies: replies, + ), + ); + + @override + Future> shareSecretShare({ + required SessionID sid, + required cl.ECCompressedPublicKey groupKey, + required Map encryptedSecrets, + }) => + _queue.run( + () => super.shareSecretShare( + sid: sid, + groupKey: groupKey, + encryptedSecrets: encryptedSecrets, + ), + ); + + @override + Future ackKeyConstructed({ + required SessionID sid, + required Signed constructedKey, + }) => + _queue.run( + () => super.ackKeyConstructed( + sid: sid, + constructedKey: constructedKey, + ), + ); +} diff --git a/pubspec.yaml b/pubspec.yaml index 5b38c39..8ae2139 100644 --- a/pubspec.yaml +++ b/pubspec.yaml @@ -13,6 +13,11 @@ dependencies: collection: ^1.17.1 grpc: ^4.0.1 args: ^2.6.0 + shelf: ^1.4.2 + shelf_router: ^1.1.4 + shelf_web_socket: ^2.0.1 + web_socket_channel: '>=2.0.0 <4.0.0' + logger: ^2.7.0 dev_dependencies: lints: ^6.0.0 diff --git a/test/config_test.dart b/test/config_test.dart index 7cb7153..16cada0 100644 --- a/test/config_test.dart +++ b/test/config_test.dart @@ -26,6 +26,12 @@ void yamlTest( final grpcConfig = GrpcConfig(server: serverConfig, port: 80); +String _indentYaml(String yaml) => yaml + .split('\n') + .where((line) => line.isNotEmpty) + .map((line) => ' $line') + .join('\n'); + void main() { setUpAll(loadFrosty); @@ -48,5 +54,13 @@ void main() { (yaml) => GrpcConfig.fromYaml(yaml), (config) => config.toHex(), ); + + test("defaults to port 50051 when YAML port is omitted", () { + final config = GrpcConfig.fromYaml( + 'server:\n${_indentYaml(serverConfig.yaml)}\n', + ); + + expect(config.port, GrpcConfig.defaultPort); + }); }); } diff --git a/test/rest_test.dart b/test/rest_test.dart new file mode 100644 index 0000000..730129f --- /dev/null +++ b/test/rest_test.dart @@ -0,0 +1,257 @@ +import 'dart:async'; +import 'dart:convert'; +import 'dart:io'; +import 'dart:typed_data'; +import 'package:noosphere_roast_server/noosphere_roast_server.dart'; +import 'package:noosphere_roast_server/src/server/state/client_session.dart'; +import 'package:noosphere_roast_server/src/server/state/state.dart'; +import 'package:shelf/shelf.dart'; +import 'package:shelf/shelf_io.dart' as shelf_io; +import 'package:test/test.dart'; + +String _b64(List bytes) => base64Encode(bytes); +Uint8List _dataBytes(String body) { + final json = jsonDecode(body) as Map; + return base64Decode(json['data'] as String); +} + +Request _jsonPost(String path, Map body) => Request( + 'POST', + Uri.parse('http://localhost$path'), + body: jsonEncode(body), + headers: {'content-type': 'application/json'}, + ); + +Request _get(String path) => Request('GET', Uri.parse('http://localhost$path')); + +Request _options(String path) => Request( + 'OPTIONS', + Uri.parse('http://localhost$path'), + ); + +Future _post( + Handler handler, + String path, + Map body, +) async => + await handler(_jsonPost(path, body)); + +Future _serve(Handler handler) => + shelf_io.serve(handler, 'localhost', 0); + +String _wsUrl(HttpServer server, String path) => + 'ws://localhost:${server.port}$path'; + +SessionID _sid([int lastByte = 1]) => + SessionID.fromBytes(Uint8List(16)..last = lastByte); + +class _FakeSession implements ClientSession { + @override + final SessionID sessionID; + + @override + final Expiry expiry; + + @override + final StreamController eventController; + + _FakeSession({ + SessionID? sid, + Expiry? expiry, + void Function()? onCancel, + }) : sessionID = sid ?? _sid(), + expiry = expiry ?? Expiry(Duration(minutes: 5)), + eventController = StreamController(onCancel: onCancel); + + void send(Event event) => sendEvent(event); + + @override + void sendEvent(Event event) => eventController.add(event); + + @override + dynamic noSuchMethod(Invocation invocation) => super.noSuchMethod(invocation); +} + +class _RestTestApi implements ServerApiHandler { + final expiry = Expiry(Duration(minutes: 5)); + final sessions = {}; + + @override + final logger = createNoosphereRoastServerLogger(); + + @override + Future extendSession(SessionID sid) async { + if (!sessions.containsKey(sid)) throw InvalidRequest.noSession(); + return expiry; + } + + @override + ClientSession getSession(SessionID id) { + final session = sessions[id]; + if (session == null) throw InvalidRequest.noSession(); + return session; + } + + @override + dynamic noSuchMethod(Invocation invocation) => super.noSuchMethod(invocation); +} + +void main() { + group('RestWebSocketNoosphereService', () { + late _RestTestApi api; + late Handler handler; + + setUp(() { + api = _RestTestApi(); + handler = RestWebSocketNoosphereService( + api: api, + allowOrigin: 'https://app.example', + ).handler; + }); + + test('handles CORS preflight requests', () async { + final response = await handler(_options('/extend-session')); + + expect(response.statusCode, 200); + expect( + response.headers['access-control-allow-origin'], + 'https://app.example', + ); + expect( + response.headers['access-control-allow-methods'], + contains('POST'), + ); + }); + + test('can leave CORS headers to a reverse proxy', () async { + final noCorsHandler = RestWebSocketNoosphereService( + api: api, + allowOrigin: null, + ).handler; + + final response = await _post(noCorsHandler, '/extend-session', {}); + + expect(response.statusCode, 400); + expect(response.headers, isNot(contains('access-control-allow-origin'))); + }); + + test('maps invalid requests to JSON errors with CORS headers', () async { + final response = await _post(handler, '/extend-session', {}); + + expect(response.statusCode, 400); + expect(response.headers['content-type'], 'application/json'); + expect( + response.headers['access-control-allow-origin'], + 'https://app.example', + ); + + final body = jsonDecode(await response.readAsString()); + expect(body, {'error': 'Expected "sid" to be a string'}); + }); + + test('extends a session through REST', () async { + final sid = _sid(); + api.sessions[sid] = _FakeSession(); + + final response = await _post(handler, '/extend-session', { + 'sid': _b64(sid.toBytes()), + }); + + expect(response.statusCode, 200); + expect( + Expiry.fromBytes(_dataBytes(await response.readAsString())) + .time + .millisecondsSinceEpoch, + api.expiry.time.millisecondsSinceEpoch, + ); + }); + + test('streams websocket events and cancels the session stream', () async { + final canceled = Completer(); + final sid = _sid(); + final session = _FakeSession( + onCancel: () { + if (!canceled.isCompleted) canceled.complete(); + }, + ); + api.sessions[sid] = session; + + final server = await _serve(handler); + try { + final socket = await WebSocket.connect( + _wsUrl(server, restWebSocketSessionPath(sid)), + ); + + session.send(KeepaliveEvent()); + + final message = await socket.first.timeout(Duration(seconds: 2)); + expect(jsonDecode(message as String), { + 'type': 'keepalive', + 'data': '', + }); + + await socket.close(); + await canceled.future.timeout(Duration(seconds: 2)); + } finally { + await server.close(force: true); + } + }); + + test('streams websocket events sent through server state fanout', () async { + final state = ServerState(logger: api.logger); + final creatorSid = _sid(1); + final receiverSid = _sid(2); + + state.clientSessions[creatorSid] = _FakeSession(sid: creatorSid); + final receiverSession = _FakeSession(sid: receiverSid); + state.clientSessions[receiverSid] = receiverSession; + api.sessions[receiverSid] = receiverSession; + + final server = await _serve(handler); + try { + final socket = await WebSocket.connect( + _wsUrl(server, restWebSocketSessionPath(receiverSid)), + ); + + final event = KeepaliveEvent(); + state.sendEventToOthers(event, creatorSid); + + final message = await socket.first.timeout(Duration(seconds: 2)); + expect(jsonDecode(message as String), { + 'type': 'keepalive', + 'data': '', + }); + + await socket.close(); + } finally { + await server.close(force: true); + } + }); + + test('returns a clean error for an unknown websocket session', () async { + final response = await handler(_get(restWebSocketSessionPath(_sid(2)))); + + expect(response.statusCode, 400); + final body = jsonDecode(await response.readAsString()); + expect(body, {'error': InvalidRequest.noSession().message}); + }); + + test('rejects websocket connections from a different origin', () async { + final sid = _sid(); + api.sessions[sid] = _FakeSession(); + + final server = await _serve(handler); + try { + await expectLater( + WebSocket.connect( + _wsUrl(server, restWebSocketSessionPath(sid)), + headers: {'Origin': 'https://other.example'}, + ), + throwsA(isA()), + ); + } finally { + await server.close(force: true); + } + }); + }); +}