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}');
}