1
0
mirror of https://github.com/danog/MadelineProto.git synced 2025-01-23 00:51:12 +01:00

Add VoIP writing functionality

This commit is contained in:
Daniil Gentili 2024-01-11 20:47:39 +01:00
parent 7e03ce81de
commit 30168b5ea5
5 changed files with 84 additions and 158 deletions

View File

@ -44,72 +44,6 @@ use function count;
*/
final class Ogg
{
private const CRC = [
0x00000000,0x04c11db7,0x09823b6e,0x0d4326d9,
0x130476dc,0x17c56b6b,0x1a864db2,0x1e475005,
0x2608edb8,0x22c9f00f,0x2f8ad6d6,0x2b4bcb61,
0x350c9b64,0x31cd86d3,0x3c8ea00a,0x384fbdbd,
0x4c11db70,0x48d0c6c7,0x4593e01e,0x4152fda9,
0x5f15adac,0x5bd4b01b,0x569796c2,0x52568b75,
0x6a1936c8,0x6ed82b7f,0x639b0da6,0x675a1011,
0x791d4014,0x7ddc5da3,0x709f7b7a,0x745e66cd,
0x9823b6e0,0x9ce2ab57,0x91a18d8e,0x95609039,
0x8b27c03c,0x8fe6dd8b,0x82a5fb52,0x8664e6e5,
0xbe2b5b58,0xbaea46ef,0xb7a96036,0xb3687d81,
0xad2f2d84,0xa9ee3033,0xa4ad16ea,0xa06c0b5d,
0xd4326d90,0xd0f37027,0xddb056fe,0xd9714b49,
0xc7361b4c,0xc3f706fb,0xceb42022,0xca753d95,
0xf23a8028,0xf6fb9d9f,0xfbb8bb46,0xff79a6f1,
0xe13ef6f4,0xe5ffeb43,0xe8bccd9a,0xec7dd02d,
0x34867077,0x30476dc0,0x3d044b19,0x39c556ae,
0x278206ab,0x23431b1c,0x2e003dc5,0x2ac12072,
0x128e9dcf,0x164f8078,0x1b0ca6a1,0x1fcdbb16,
0x018aeb13,0x054bf6a4,0x0808d07d,0x0cc9cdca,
0x7897ab07,0x7c56b6b0,0x71159069,0x75d48dde,
0x6b93dddb,0x6f52c06c,0x6211e6b5,0x66d0fb02,
0x5e9f46bf,0x5a5e5b08,0x571d7dd1,0x53dc6066,
0x4d9b3063,0x495a2dd4,0x44190b0d,0x40d816ba,
0xaca5c697,0xa864db20,0xa527fdf9,0xa1e6e04e,
0xbfa1b04b,0xbb60adfc,0xb6238b25,0xb2e29692,
0x8aad2b2f,0x8e6c3698,0x832f1041,0x87ee0df6,
0x99a95df3,0x9d684044,0x902b669d,0x94ea7b2a,
0xe0b41de7,0xe4750050,0xe9362689,0xedf73b3e,
0xf3b06b3b,0xf771768c,0xfa325055,0xfef34de2,
0xc6bcf05f,0xc27dede8,0xcf3ecb31,0xcbffd686,
0xd5b88683,0xd1799b34,0xdc3abded,0xd8fba05a,
0x690ce0ee,0x6dcdfd59,0x608edb80,0x644fc637,
0x7a089632,0x7ec98b85,0x738aad5c,0x774bb0eb,
0x4f040d56,0x4bc510e1,0x46863638,0x42472b8f,
0x5c007b8a,0x58c1663d,0x558240e4,0x51435d53,
0x251d3b9e,0x21dc2629,0x2c9f00f0,0x285e1d47,
0x36194d42,0x32d850f5,0x3f9b762c,0x3b5a6b9b,
0x0315d626,0x07d4cb91,0x0a97ed48,0x0e56f0ff,
0x1011a0fa,0x14d0bd4d,0x19939b94,0x1d528623,
0xf12f560e,0xf5ee4bb9,0xf8ad6d60,0xfc6c70d7,
0xe22b20d2,0xe6ea3d65,0xeba91bbc,0xef68060b,
0xd727bbb6,0xd3e6a601,0xdea580d8,0xda649d6f,
0xc423cd6a,0xc0e2d0dd,0xcda1f604,0xc960ebb3,
0xbd3e8d7e,0xb9ff90c9,0xb4bcb610,0xb07daba7,
0xae3afba2,0xaafbe615,0xa7b8c0cc,0xa379dd7b,
0x9b3660c6,0x9ff77d71,0x92b45ba8,0x9675461f,
0x8832161a,0x8cf30bad,0x81b02d74,0x857130c3,
0x5d8a9099,0x594b8d2e,0x5408abf7,0x50c9b640,
0x4e8ee645,0x4a4ffbf2,0x470cdd2b,0x43cdc09c,
0x7b827d21,0x7f436096,0x7200464f,0x76c15bf8,
0x68860bfd,0x6c47164a,0x61043093,0x65c52d24,
0x119b4be9,0x155a565e,0x18197087,0x1cd86d30,
0x029f3d35,0x065e2082,0x0b1d065b,0x0fdc1bec,
0x3793a651,0x3352bbe6,0x3e119d3f,0x3ad08088,
0x2497d08d,0x2056cd3a,0x2d15ebe3,0x29d4f654,
0xc5a92679,0xc1683bce,0xcc2b1d17,0xc8ea00a0,
0xd6ad50a5,0xd26c4d12,0xdf2f6bcb,0xdbee767c,
0xe3a1cbc1,0xe760d676,0xea23f0af,0xeee2ed18,
0xf0a5bd1d,0xf464a0aa,0xf9278673,0xfde69bc4,
0x89b8fd09,0x8d79e0be,0x803ac667,0x84fbdbd0,
0x9abc8bd5,0x9e7d9662,0x933eb0bb,0x97ffad0c,
0xafb010b1,0xab710d06,0xa6322bdf,0xa2f33668,
0xbcb4666d,0xb8757bda,0xb5365d03,0xb1f740b4,
];
private const OPUS_SET_APPLICATION_REQUEST = 4000;
private const OPUS_GET_APPLICATION_REQUEST = 4001;
private const OPUS_SET_BITRATE_REQUEST = 4002;
@ -657,80 +591,16 @@ final class Ogg
$checkErr($opus->opus_encoder_ctl($encoder, self::OPUS_SET_BANDWIDTH_REQUEST, self::OPUS_BANDWIDTH_FULLBAND));
$checkErr($opus->opus_encoder_ctl($encoder, self::OPUS_SET_BITRATE_REQUEST, 130*1000));
$out = $oggOut instanceof LocalFile
? openFile($oggOut->file, 'w')
: $oggOut;
$writePage = static function (int $header_type_flag, int $granule, int $streamId, int &$streamSeqno, string $packet) use ($out): void {
Assert::true(\strlen($packet) < 65025);
$segments = [
...array_fill(0, (int) (\strlen($packet) / 255), 255),
\strlen($packet) % 255,
];
$data = 'OggS'.pack(
'CCPVVVCC*',
0, // stream_structure_version
$header_type_flag,
$granule,
$streamId,
$streamSeqno++,
0,
\count($segments),
...$segments
).$packet;
$c = 0;
for ($i = 0; $i < \strlen($data); $i++) {
$c = ($c<<8)^self::CRC[(($c >> 24)&0xFF)^(\ord($data[$i]))];
}
$crc = pack('V', $c);
$data = substr_replace(
$data,
$crc,
22,
4
);
$out->write($data);
};
$streamId = unpack('V', Tools::random(4))[1];
$seqno = 0;
$writePage(
Ogg::BOS,
0,
$streamId,
$seqno,
'OpusHead'.pack(
'CCvVvC',
1,
$header['channels'],
312,
$header['sampleRate'],
0,
0,
)
if ($oggOut instanceof LocalFile) {
$oggOut = openFile($oggOut->file, 'w');
}
$writer = new OggWriter($oggOut);
$writer->writeHeader(
$header['channels'],
$header['sampleRate'],
$opus->opus_get_version_string()
);
$tags = 'OpusTags';
$writeTag = static function (string $tag) use (&$tags): void {
$tags .= pack('V', \strlen($tag)).$tag;
};
$writeTag("MadelineProto ".API::RELEASE.", ".$opus->opus_get_version_string());
$tags .= pack('V', 2);
$writeTag("ENCODER=MadelineProto ".API::RELEASE." with ".$opus->opus_get_version_string());
$writeTag("MADELINE_ENCODER_V=1");
$writeTag('See https://docs.madelineproto.xyz/docs/CALLS.html for more info');
$writePage(
0,
0,
$streamId,
$seqno,
$tags
);
$granule = 0;
$buf = $opus->cast($opus->type('char*'), FFI::addr($opus->new('char[1024]')));
do {
$chunkOrig = $read($chunkSize) ?? '';
@ -738,17 +608,13 @@ final class Ogg
$granuleDiff = \strlen($chunk) >> $shift;
$len = $opus->opus_encode($encoder, $chunk, $granuleDiff, $buf, 1024);
$checkErr($len);
$writePage(
\strlen($chunk) !== \strlen($chunkOrig) ? self::EOS : 0,
$granule += $granuleDiff,
$streamId,
$seqno,
FFI::string($buf, $len)
$writer->writeChunk(
FFI::string($buf, $len),
$granuleDiff,
\strlen($chunk) !== \strlen($chunkOrig)
);
} while (\strlen($chunk) === \strlen($chunkOrig));
$opus->opus_encoder_destroy($encoder);
unset($buf, $encoder, $opus);
$out->close();
}
}

View File

@ -17,6 +17,7 @@
namespace danog\MadelineProto;
use Amp\ByteStream\ReadableStream;
use Amp\ByteStream\WritableStream;
use danog\MadelineProto\EventHandler\SimpleFilters;
use danog\MadelineProto\EventHandler\Update;
use danog\MadelineProto\VoIP\CallState;
@ -98,6 +99,18 @@ final class VoIP extends Update implements SimpleFilters
return $this;
}
/**
* Set output file or stream for incoming OPUS audio packets.
*
* Will write an OGG OPUS stream to the specified file or stream.
*/
public function setOutput(LocalFile|WritableStream $file): self
{
$this->getClient()->callSetOutput($this->callID, $file);
return $this;
}
/**
* Play file.
*/

View File

@ -21,6 +21,7 @@ declare(strict_types=1);
namespace danog\MadelineProto\VoIP;
use Amp\ByteStream\ReadableStream;
use Amp\ByteStream\WritableStream;
use Amp\DeferredFuture;
use AssertionError;
use danog\MadelineProto\LocalFile;
@ -177,6 +178,16 @@ trait AuthKeyHandler
($this->calls[$id] ?? null)?->play($file);
}
/**
* Set output file or stream for incoming OPUS audio packets in a call.
*
* Will write an OGG OPUS stream to the specified file or stream.
*/
public function callSetOutput(int $id, LocalFile|WritableStream $file): void
{
($this->calls[$id] ?? null)?->setOutput($file);
}
/**
* Play file in call, blocking until the file has finished playing if a stream is provided.
*

View File

@ -324,17 +324,6 @@ final class Endpoint
$result['id'] = \ord(stream_get_contents($message, 1));
$result['enabled'] = \ord(stream_get_contents($message, 1));
break;
// streamData flags:int2 stream_id:int6 has_more_flags:flags.1?true length:(flags.0?int16:int8) timestamp:int data:byteArray = StreamData;
// packetStreamData#4 stream_data:streamData = Packet;
case VoIPController::PKT_STREAM_DATA:
$flags = \ord(stream_get_contents($message, 1));
$result['stream_id'] = $flags & 0x3F;
$flags = ($flags & 0xC0) >> 6;
$result['has_more_flags'] = (bool) ($flags & 2);
$length = $flags & 1 ? unpack('v', stream_get_contents($message, 2))[1] : \ord(stream_get_contents($message, 1));
$result['timestamp'] = unpack('V', stream_get_contents($message, 4))[1];
$result['data'] = stream_get_contents($message, $length);
break;
case \danog\MadelineProto\VoIPController::PKT_UPDATE_STREAMS:
continue 2;
case \danog\MadelineProto\VoIPController::PKT_PING:
@ -345,6 +334,17 @@ final class Endpoint
$result['out_seq_no'] = unpack('V', stream_get_contents($payload, 4))[1];
}
break;
// streamData flags:int2 stream_id:int6 has_more_flags:flags.1?true length:(flags.0?int16:int8) timestamp:int data:byteArray = StreamData;
// packetStreamData#4 stream_data:streamData = Packet;
case VoIPController::PKT_STREAM_DATA:
$flags = \ord(stream_get_contents($message, 1));
$result[0]['stream_id'] = $flags & 0x3F;
$flags = ($flags & 0xC0) >> 6;
$result[0]['has_more_flags'] = (bool) ($flags & 2);
$length = $flags & 1 ? unpack('v', stream_get_contents($message, 2))[1] : \ord(stream_get_contents($message, 1));
$result[0]['timestamp'] = unpack('V', stream_get_contents($message, 4))[1];
$result[0]['data'] = stream_get_contents($message, $length);
break;
case VoIPController::PKT_STREAM_DATA_X2:
for ($x = 0; $x < 2; $x++) {
$flags = \ord(stream_get_contents($message, 1));

View File

@ -21,6 +21,7 @@ namespace danog\MadelineProto;
use Amp\ByteStream\BufferedReader;
use Amp\ByteStream\ReadableBuffer;
use Amp\ByteStream\ReadableStream;
use Amp\ByteStream\WritableStream;
use Amp\Sync\LocalMutex;
use danog\Loop\Loop;
use danog\MadelineProto\Loop\VoIP\DjLoop;
@ -37,6 +38,7 @@ use Webmozart\Assert\Assert;
use function Amp\async;
use function Amp\delay;
use function Amp\File\openFile;
use function Amp\Future\await;
/** @internal */
@ -144,6 +146,10 @@ final class VoIPController
private LocalMutex $authMutex;
private ?LocalFile $outputFile = null;
private ?OggWriter $outputStream = null;
private ?int $outputStreamId = null;
/**
* Constructor.
*
@ -169,7 +175,8 @@ final class VoIPController
public function __serialize(): array
{
$result = get_object_vars($this);
unset($result['authMutex']);
unset($result['authMutex'], $result['outputStream']);
return $result;
}
/**
@ -187,6 +194,11 @@ final class VoIPController
$this->diskJockey ??= new DjLoop($this);
Assert::true($this->diskJockey->start());
EventLoop::queue(function (): void {
if (isset($this->outputFile)) {
\assert($this->outputStreamId !== null);
$out = openFile($this->outputFile->file, 'a');
$this->outputStream = new OggWriter($out, $this->outputStreamId);
}
if ($this->callState === CallState::RUNNING) {
if ($this->voipState === VoIPState::CREATED) {
// No endpoints yet
@ -588,6 +600,7 @@ final class VoIPController
*/
private function handlePacket(Endpoint $socket, array $packet): void
{
$cnt = 0;
switch ($packet['_']) {
case self::PKT_INIT:
$this->setVoipState(VoIPState::WAIT_INIT_ACK);
@ -628,6 +641,11 @@ final class VoIPController
}
break;
}
if ($this->outputStream !== null && $cnt) {
foreach ($packet as ['data' => $data]) {
$this->outputStream->writeChunk($data, 2880, false);
}
}
}
private function initStream(): void
{
@ -750,6 +768,24 @@ final class VoIPController
$t += $delay;
}
}
/**
* Set output file or stream for incoming OPUS audio packets.
*
* Will write an OGG OPUS stream to the specified file or stream.
*/
public function setOutput(LocalFile|WritableStream $file): void
{
if ($file instanceof LocalFile) {
$this->outputFile = $file;
$file = openFile($file->file, 'w');
} else {
$this->outputFile = null;
}
$this->outputStreamId = random_int(-(2**31), (2**31)-1);
$this->outputStream = new OggWriter($file, $this->outputStreamId);
$this->outputStream->writeHeader(1, 48000, "incoming audio stream");
}
/**
* Play file.
*/