start method

Future<void> start()

Implementation

Future<void> start() async {
  if (_isRunning) return;

  if (!_routesRegistered) {
    _router.private.post('/api/v1/events/publish', (Request request) async {
      try {
        final payload = await request.readAsString();
        final json = jsonDecode(payload);
        final event = EventModel.fromJson(json);

        if (await _db.events.id(event.hash).get() != null) {
          return Response.ok('Already Received');
        }

        if (await event.verify()) {
          late final bool ingested;
          try {
            ingested = await EventIngestor.instance.ingest(event);
          } catch (_) {
            final storedEvent = await _db.events.id(event.hash).get();
            if (storedEvent != null) {
              return Response.ok('Received');
            }
            rethrow;
          }

          final storedEvent = await _db.events.id(event.hash).get();

          if (storedEvent == null) {
            return PeerResponse.invalidPayload(
              message: 'Event was not accepted by protocol validation',
            );
          }

          if (ingested) {
            await _db.logs.insert(
              EventLogModel(
                type: EventLogTypes.publishReceived,
                deviceKey: request.device.signingKey,
                body: PublishReceivedLogBodyModel(
                  eventHash: event.hash,
                  deviceKey: event.deviceKey,
                ),
              ),
            );

            client.events.relay(event);
          }

          return Response.ok('Received');
        } else {
          return PeerResponse.invalidSignature();
        }
      } catch (error) {
        if (error is FormatException || error is TypeError) rethrow;
        return PeerResponse.internalError(message: 'Error processing event');
      }
    });

    _router.private.post('/api/v1/events/sync', (Request request) async {
      try {
        final payload = await request.readAsString();
        final json = jsonDecode(payload) as Map<String, dynamic>;
        final index = json.map((key, value) => MapEntry(key, value as int));

        final missingEvents = await _db.events.diff(index, limit: 1000).get();

        final responseBody = jsonEncode({
          'events': missingEvents.map((event) => event.toJson()).toList(),
          'has_more': missingEvents.length == 1000,
        });

        if (missingEvents.isNotEmpty) {
          await _db.logs.insert(
            EventLogModel(
              type: EventLogTypes.syncSent,
              deviceKey: request.device.signingKey,
              body: SyncSentLogBodyModel(
                deviceKey: request.device.signingKey,
                count: missingEvents.length,
              ),
            ),
          );
        }

        return Response.ok(
          responseBody,
          headers: {'Content-Type': 'application/json'},
        );
      } catch (error) {
        if (error is FormatException || error is TypeError) rethrow;
        logger.log('Error processing sync: $error');
        return PeerResponse.internalError(message: 'Error processing sync');
      }
    });

    _router.private.post('/api/v1/canvas/sync', (Request request) async {
      try {
        final payload = await request.readAsString();
        final json = jsonDecode(payload);
        if (json is! Map<String, dynamic>) {
          return PeerResponse.invalidPayload(
            message: 'Invalid canvas sync payload',
          );
        }

        late final CanvasSyncRequestModel syncRequest;
        try {
          syncRequest = CanvasSyncRequestModel.fromJson(json);
        } catch (_) {
          return PeerResponse.invalidPayload(
            message: 'Invalid canvas sync payload',
          );
        }

        final requesterDevice = request.device;
        final conversationKey = syncRequest.conversationKey;
        final canvasId = syncRequest.canvasId;

        final membership = await _db.conversationParticipants
            .id(conversationKey, requesterDevice.peerKey)
            .get();
        if (membership == null) {
          return PeerResponse.notParticipant();
        }

        final canvas = await _db.canvases.id(canvasId).get();
        if (canvas == null || canvas.conversationKey != conversationKey) {
          return PeerResponse.notFound(message: 'Canvas not found');
        }

        final objectIds = await _db.canvasStateSync.changedObjectIds(
          conversationKey,
          canvasId: canvasId,
          observedActorHlcWatermarks: syncRequest.objectActorHlcWatermarks,
          limit: _kCanvasSyncMaxObjects + 1,
        );

        if (objectIds.length > _kCanvasSyncMaxObjects) {
          final error = CanvasSyncErrorModel(
            code: CanvasSyncErrorCode.syncTooLarge,
            conversationKey: syncRequest.conversationKey,
            canvasId: canvasId,
            objectCount: objectIds.length,
            maxObjects: _kCanvasSyncMaxObjects,
            maxPayloadBytes: _kCanvasSyncMaxPayloadBytes,
            message: 'Canvas sync exceeds object cap',
          );
          return Response(
            413,
            body: jsonEncode(error.toJson()),
            headers: {'Content-Type': 'application/json'},
          );
        }

        final responseModel = await _db.canvasStateSync.snapshot(
          conversationKey,
          canvasId: canvasId,
          objectIds: objectIds,
          observedMetadataActorHlcWatermarks:
              syncRequest.metadataActorHlcWatermarks,
        );

        final responseBody = jsonEncode(responseModel.toJson());
        final payloadBytes = utf8.encode(responseBody).length;
        if (payloadBytes > _kCanvasSyncMaxPayloadBytes) {
          final error = CanvasSyncErrorModel(
            code: CanvasSyncErrorCode.syncTooLarge,
            conversationKey: syncRequest.conversationKey,
            canvasId: canvasId,
            objectCount: responseModel.objects.length,
            maxObjects: _kCanvasSyncMaxObjects,
            payloadBytes: payloadBytes,
            maxPayloadBytes: _kCanvasSyncMaxPayloadBytes,
            message: 'Canvas sync exceeds payload byte cap',
          );
          return Response(
            413,
            body: jsonEncode(error.toJson()),
            headers: {'Content-Type': 'application/json'},
          );
        }

        return Response.ok(
          responseBody,
          headers: {'Content-Type': 'application/json'},
        );
      } catch (error) {
        logger.log('Error processing canvas sync: $error');
        return PeerResponse.internalError(
          message: 'Error processing canvas sync',
        );
      }
    });

    _router.private.get('/api/v1/media/<hash>', (
      Request request,
      String hash,
    ) async {
      if (!MediaStorage.isValidHash(hash)) {
        return PeerResponse.invalidPayload(message: 'Invalid media hash');
      }

      final bytes = await MediaStorage.instance.get(hash);
      if (bytes == null) {
        return PeerResponse.notFound(message: 'Media not found');
      }

      return Response.ok(
        bytes,
        headers: {
          'Content-Type': 'application/octet-stream',
          'Content-Length': bytes.length.toString(),
        },
      );
    });

    /// Device identification is a public endpoint used to communicate device identity between devices.
    _router.public.get('/api/v1/identity', (Request request) async {
      try {
        final clientIp = request.requiredClientIp;
        final myIp = await client.getCurrentIp().required();
        final binding = await DeviceRequestBindingModel.create(
          clientIp: clientIp,
          serverIp: myIp,
          method: 'GET',
          path: '/api/v1/identity',
        );

        final authEvent = await _db.events
            .deviceAuthEvent(binding.deviceKey)
            .get()
            .required();

        return Response.ok(
          jsonEncode({
            'binding': binding.toJson(),
            'auth_event': authEvent.toJson(),
          }),
          headers: {'Content-Type': 'application/json'},
        );
      } catch (error) {
        logger.log('Error processing identity handshake: $error');
        return PeerResponse.internalError(
          message: 'Error processing identity',
        );
      }
    });

    _router.private.post('/api/v1/admin/auth-key', (Request request) async {
      try {
        final tailnetId = TailscaleAdmin.instance.tailnetId;
        final apiKey = TailscaleAdmin.instance.adminApiKey;

        if (tailnetId == null || apiKey == null) {
          return PeerResponse.notAdmin();
        }

        final deviceKey = request.device.signingKey;

        final twentyFourHoursAgo = DateTime.now().subtract(
          const Duration(hours: 24),
        );
        final requestCount = await _db.logs
            .authKeyRequestCount(deviceKey, twentyFourHoursAgo)
            .get();

        if (requestCount >= 5) {
          return PeerResponse.rateLimitExceeded();
        }

        final authKey = await TailscaleAdmin.instance.generateAuthKey(
          apiKey: apiKey,
          tailnetId: tailnetId,
        );

        await _db.logs.insert(
          EventLogModel(
            type: EventLogTypes.authKeyRequest,
            deviceKey: deviceKey,
            body: AuthKeyRequestLogBodyModel(deviceKey: deviceKey),
          ),
        );

        return Response.ok(
          authKey,
          headers: {'Content-Type': 'application/json'},
        );
      } catch (error) {
        logger.log('Error processing key request: $error');
        return PeerResponse.internalError(
          message: 'Error processing key request',
        );
      }
    });

    _router.private.post('/api/v1/messages/grants', (Request request) async {
      try {
        final device = request.device;
        if (device.peerKey != client.peerPublicKey) {
          return PeerResponse.peerMismatch();
        }

        final currentEncryptionKey = Encryption.instance.publicKey;
        final requesterSigningKey = device.signingKey;
        final requesterEncryptionKey = device.encryptionKey;

        final messageEvents = await _db.events.currentMessageEvents.get();

        final grants = (await Future.wait(
          messageEvents.map((event) async {
            final body = event.body;

            if (body.keys.containsKey(requesterSigningKey)) {
              return null;
            }

            final secretKey = await Messenger.instance.openSecretKey(event);
            if (secretKey == null) {
              logger.log('Unable to get secret key for ${event.hash}');
              return null;
            }

            final reencryptedKey = base64Encode(
              await _encryption.encryptSharedSecret(
                await secretKey.extractBytes(),
                requesterEncryptionKey,
              ),
            );

            return MessageGrantModel(
              hash: event.hash,
              issuerEncryptionKey: currentEncryptionKey,
              encryptedSecretKey: reencryptedKey,
            );
          }),
        )).whereType<MessageGrantModel>().toList();

        return Response.ok(
          jsonEncode({
            'grants': grants.map((grant) => grant.toJson()).toList(),
          }),
          headers: {'Content-Type': 'application/json'},
        );
      } catch (error) {
        logger.log('Error processing grants: $error');
        return PeerResponse.internalError(message: 'Error processing grants');
      }
    });

    _router.private.post('/api/v1/messages/receipts', (
      Request request,
    ) async {
      try {
        final device = request.device;
        final payload = await request.readAsString();

        final ConversationReceiptsRequestModel(
          :conversationKey,
          :startingAt,
          :updatedAfter,
        ) = ConversationReceiptsRequestModel.fromJson(
          jsonDecode(payload),
        );

        final participant = await _db.conversationParticipants
            .id(conversationKey, device.peerKey)
            .get();
        if (participant == null) {
          return PeerResponse.notParticipant();
        }

        final timestamp = DateTime.now();

        final receipts = await _db.messageReceipts
            .conversation(
              conversationKey,
              startingAt: startingAt,
              updatedAfter: updatedAfter,
            )
            .get();

        final response = ConversationReceiptsResponseModel(
          receipts: receipts,
          timestamp: timestamp,
        );

        return Response.ok(
          jsonEncode(response.toJson()),
          headers: {'Content-Type': 'application/json'},
        );
      } catch (error) {
        logger.log('Error processing receipt sync: $error');
        return PeerResponse.internalError(
          message: 'Error processing receipt sync',
        );
      }
    });
    _routesRegistered = true;
  }

  final publicPipeline = const Pipeline()
      .addMiddleware(middleware.clientIpHeaderMiddleware())
      .addMiddleware(middleware.errorBoundaryMiddleware())
      .addMiddleware(middleware.loopbackMiddleware());

  final privatePipeline = publicPipeline
      .addMiddleware(middleware.bodySizeLimitMiddleware())
      .addMiddleware(middleware.deviceBindingMiddleware())
      .addMiddleware(middleware.rateLimitMiddleware());

  final handler = Cascade()
      .add(publicPipeline.addHandler(_router.public.call))
      .add(privatePipeline.addHandler(_router.private.call))
      .handler;

  _server = await shelf_io.serve(handler, InternetAddress.loopbackIPv4, 0);
  _isRunning = true;
  logger.log('P2P Server listening on port ${_server.port}');
}