/
check-sync-on-feeds
executable file
·280 lines (253 loc) · 10.1 KB
/
check-sync-on-feeds
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
#!/usr/bin/python
# -*- Mode: Python -*-
# vi:si:et:sw=4:sts=4:ts=4
# (C) Copyright 2007 Zaheer Abbas Merali <zaheerabbas at merali dot org>
#
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation; either version 2 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU Library General Public License for more details.
#
# You should have received a copy of the GNU General Public License
# along with this program; if not, write to the Free Software
# Foundation, Inc., 59 Temple Place - Suite 330, Boston, MA 02111-1307, USA.
from flumotion.component import feed
from twisted.internet import reactor
from flumotion.twisted import pb
from flumotion.common import log, errors
from flumotion.admin import connections
from flumotion.admin.command import utils
from flumotion.admin.admin import AdminModel
import os
import sys
import string
import gobject
gobject.threads_init()
import gst
import optparse
haveVideoFeed = False
haveAudioFeed = False
videoFd = 0
audioFd = 0
videoBufferProbeId = 0
audioBufferProbeId = 0
haveVideoBuffer = False
haveAudioBuffer = False
lowestTimestamp = 0L
pipe = None
sentNewsegment = False
seenNewsegmentAudio = False
seenNewsegmentVideo = False
def usage(args, exitval=0):
print "usage: %s [OPTIONS] -m MANAGER " \
"-V FULLFEEDIDOFVIDEOFEEDER -A FULLFEEDIDOFAUDIOFEEDER" % args[0]
print ''
print 'See %s -h for help on the available options.' % args[0]
sys.exit(exitval)
def gotVideoFeed(res):
global haveAudioFeed
global haveVideoFeed
global videoFd
if not res:
log.debug("output-feed", "got None in gotFeed")
reactor.stop()
return
(feedId, fd) = res
videoFd = fd
log.debug("check-sync-on-feeds","Got feed on fd %r for feedId %s" %
(fd, feedId))
haveVideoFeed = True
if haveAudioFeed:
startPipeline()
def videoEventProbe(pad, event):
global seenNewsegmentVideo
if event.type == gst.EVENT_NEWSEGMENT:
log.debug("check-sync-on-feeds", "New segment event on video pad %r",
event)
if not seenNewsegmentVideo:
seenNewsegmentVideo = True
return False
return True
def audioEventProbe(pad, event):
global seenNewsegmentAudio
if event.type == gst.EVENT_NEWSEGMENT:
log.debug("check-sync-on-feeds", "New segment event on audio pad %r",
event)
if not seenNewsegmentAudio:
seenNewsegmentAudio = True
return False
return True
def sendNewsegments():
global pipe
global lowestTimestamp
global sentNewsegment
global videoBufferProbeId
global audioBufferProbeId
if not sentNewsegment:
sentNewsegment = True
log.debug("check-sync-on-feeds",
"new segment created with timestamp %d (%s)",
lowestTimestamp, gst.TIME_ARGS(lowestTimestamp))
newseg = gst.event_new_new_segment(False, 1.0, gst.FORMAT_TIME, lowestTimestamp, -1, 0)
vsinkpad = pipe.get_by_name("vsinkpad").get_pad("sink")
asinkpad = pipe.get_by_name("asinkpad").get_pad("sink")
vsinkpad.send_event(newseg)
asinkpad.send_event(newseg)
vsinkpad.remove_buffer_probe(videoBufferProbeId)
asinkpad.remove_buffer_probe(audioBufferProbeId)
def videoBufferProbe(pad, buffer):
global haveVideoBuffer
global haveAudioBuffer
global lowestTimestamp
haveVideoBuffer = True
log.debug("check-sync-on-feeds",
"video buffer arrived with timestamp %d (%s)",
buffer.timestamp, gst.TIME_ARGS(buffer.timestamp))
if haveAudioBuffer:
if lowestTimestamp > buffer.timestamp:
lowestTimestamp = buffer.timestamp
sendNewsegments()
else:
lowestTimestamp = buffer.timestamp
return False
def audioBufferProbe(pad, buffer):
global haveVideoBuffer
global haveAudioBuffer
global lowestTimestamp
haveAudioBuffer = True
log.debug("check-sync-on-feeds",
"audio buffer arrived with timestamp %d (%s)",
buffer.timestamp, gst.TIME_ARGS(buffer.timestamp))
if haveVideoBuffer:
if lowestTimestamp > buffer.timestamp:
lowestTimestamp = buffer.timestamp
sendNewsegments()
else:
lowestTimestamp = buffer.timestamp
return False
def startPipeline():
global videoFd
global audioFd
global videoBufferProbeId
global audioBufferProbeId
global pipe
log.debug("check-sync-on-feeds", "Starting pipeline")
pipe = gst.parse_launch("fdsrc fd=%d ! queue ! gdpdepay ! ffmpegcolorspace name=vsinkpad ! videoscale ! ximagesink fdsrc fd=%d ! queue ! gdpdepay ! audioconvert name=asinkpad ! audioresample ! alsasink" % (videoFd,audioFd))
vsinkpadElement = pipe.get_by_name("vsinkpad")
asinkpadElement = pipe.get_by_name("asinkpad")
# add event probes, when new segments received, block pad
vsinkpad = vsinkpadElement.get_pad("sink")
asinkpad = asinkpadElement.get_pad("sink")
vsinkpad.add_event_probe(videoEventProbe)
asinkpad.add_event_probe(audioEventProbe)
videoBufferProbeId = vsinkpad.add_buffer_probe(videoBufferProbe)
audioBufferProbeId = asinkpad.add_buffer_probe(audioBufferProbe)
pipe.set_state(gst.STATE_PLAYING)
def gotAudioFeed(res):
global haveAudioFeed
global haveVideoFeed
global audioFd
if not res:
log.debug("check-sync-on-feeds", "got None in gotFeed")
reactor.stop()
return
(feedId, fd) = res
audioFd = fd
log.debug("check-sync-on-feeds","Got feed on fd %r for feedId %s" %
(fd, feedId))
haveAudioFeed = True
if haveVideoFeed:
startPipeline()
def main(args):
log.init()
parser = optparse.OptionParser()
parser.add_option('-d', '--debug',
action="store", type="string", dest="debug",
help="set debug levels")
parser.add_option('-u', '--usage',
action="store_true", dest="usage",
help="show a usage message")
parser.add_option('-m', '--manager',
action="store", type="string", dest="manager",
help="the manager to connect to, e.g. localhost:7531")
parser.add_option('', '--no-ssl',
action="store_true", dest="no_ssl",
help="disable encryption when connecting to the manager")
parser.add_option('-V', '--video-feed-id',
action="store", type="string", dest="videoFeedId",
help="the full feed id of the video feed to connect to"
", e.g. /default/video-source:default")
parser.add_option('-A', '--audio-feed-id',
action="store", type="string", dest="audioFeedId",
help="the full feed id of the audio feed to connect to"
", e.g. /default/audio-source:default")
options, args = parser.parse_args(args)
if options.debug:
log.setFluDebug(options.debug)
if options.usage:
usage(args)
if not options.manager or not options.videoFeedId or not \
options.audioFeedId:
usage(args)
connection = connections.parsePBConnectionInfo(options.manager,
not options.no_ssl)
model = AdminModel()
d = model.connectToManager(connection)
def failed(failure):
if failure.check(errors.ConnectionRefusedError):
print "Manager refused connection. Check your user and password."
elif failure.check(errors.ConnectionFailedError):
message = "".join(failure.value.args)
print "Connection to manager failed: %s" % message
else:
print ("Exception while connecting to manager: %s"
% log.getFailureMessage(failure))
return failure
d.addErrback(failed)
d.addCallback(managerConnected, options)
reactor.run()
def getFeed(feedClient, worker, authenticator, feedId, model):
fsd = model.workerCallRemote(worker.get("name"), "getFeedServerPort")
def gotFeedServerPort(port):
return feedClient.requestFeed(worker.get('host'), port,
authenticator, feedId)
fsd.addCallback(gotFeedServerPort)
return fsd
def managerConnected(model, options):
psd = model.callRemote('getPlanetState')
def gotPlanetState(planet):
videoAvatarId = utils.avatarId(options.videoFeedId.split(":")[0])
audioAvatarId = utils.avatarId(options.audioFeedId.split(":")[0])
videoComponent = utils.find_component(planet, videoAvatarId)
audioComponent = utils.find_component(planet, audioAvatarId)
if videoComponent and audioComponent:
whsd = model.callRemote('getWorkerHeavenState')
def gotWorkerHeavenState(whs):
# find worker for each component
vworkername = videoComponent.get("workerRequested")
aworkername = audioComponent.get("workerRequested")
vworker = None
aworker = None
for worker in whs.get('workers'):
if worker.get('name') == vworkername:
vworker = worker
if worker.get('name') == aworkername:
aworker = worker
if vworker and aworker:
vclient = feed.FeedMedium(logName="check-sync-on-feeds")
aclient = feed.FeedMedium(logName="check-sync-on-feeds")
authenticator = model.connectionInfo.authenticator
vd = getFeed(vclient, vworker, authenticator,
options.videoFeedId, model)
ad = getFeed(aclient, aworker, authenticator,
options.audioFeedId, model)
vd.addCallback(gotVideoFeed)
ad.addCallback(gotAudioFeed)
whsd.addCallback(gotWorkerHeavenState)
psd.addCallback(gotPlanetState)
main(sys.argv)