Skip to content
This repository was archived by the owner on Jun 3, 2020. It is now read-only.

Commit 820c8ea

Browse files
committed
Add beaker.crawler API
1 parent 58acd25 commit 820c8ea

5 files changed

Lines changed: 76 additions & 1 deletion

File tree

crawler/index.js

Lines changed: 51 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,13 @@
1+
const emitStream = require('emit-stream')
2+
const {URL} = require('url')
13
const _throttle = require('lodash.throttle')
24
const lock = require('../lib/lock')
35
const db = require('../dbs/profile-data-db')
6+
const archivesDb = require('../dbs/archives')
47
const users = require('../users')
58
const dat = require('../dat')
69

10+
const {crawlerEvents} = require('./util')
711
const posts = require('./posts')
812
const followgraph = require('./followgraph')
913
const siteDescriptions = require('./site-descriptions')
@@ -21,6 +25,7 @@ const watches = {}
2125
exports.posts = posts
2226
exports.followgraph = followgraph
2327
exports.siteDescriptions = siteDescriptions
28+
const createEventsStream = exports.createEventsStream = () => emitStream(crawlerEvents)
2429

2530
exports.setup = async function () {
2631
}
@@ -32,6 +37,7 @@ exports.watchSite = async function (archive) {
3237
console.log('watchSite', archive.url)
3338

3439
if (!(archive.url in watches)) {
40+
crawlerEvents.emit('watch', {sourceUrl: archive.url})
3541
const queueCrawl = _throttle(() => crawlSite(archive), 5e3)
3642

3743
// watch for file changes
@@ -57,6 +63,7 @@ exports.watchSite = async function (archive) {
5763
exports.unwatchSite = async function (url) {
5864
// stop watching for file changes
5965
if (url in watches) {
66+
crawlerEvents.emit('unwatch', {sourceUrl: url})
6067
watches[url].close()
6168
watches[url] = null
6269
}
@@ -65,6 +72,7 @@ exports.unwatchSite = async function (url) {
6572
const crawlSite =
6673
exports.crawlSite = async function (archive) {
6774
console.log('crawling', archive.url)
75+
crawlerEvents.emit('crawl-start', {sourceUrl: archive.url})
6876
var release = await lock('crawl:' + archive.url)
6977
try {
7078
// get/create crawl source
@@ -80,7 +88,49 @@ exports.crawlSite = async function (archive) {
8088
followgraph.crawlSite(archive, crawlSource),
8189
siteDescriptions.crawlSite(archive, crawlSource)
8290
])
91+
} catch (err) {
92+
crawlerEvents.emit('crawl-error', {sourceUrl: archive.url, err: err.toString()})
8393
} finally {
94+
crawlerEvents.emit('crawl-finish', {sourceUrl: archive.url})
8495
release()
8596
}
86-
}
97+
}
98+
99+
const getCrawlStates =
100+
exports.getCrawlStates = async function () {
101+
var rows = await db.all(`
102+
SELECT
103+
crawl_sources.url AS url,
104+
GROUP_CONCAT(crawl_sources_meta.crawlSourceVersion) AS versions,
105+
GROUP_CONCAT(crawl_sources_meta.crawlDataset) AS datasets,
106+
MAX(crawl_sources_meta.updatedAt) AS updatedAt
107+
FROM crawl_sources
108+
INNER JOIN crawl_sources_meta ON crawl_sources_meta.crawlSourceId = crawl_sources.id
109+
GROUP BY crawl_sources.id
110+
`)
111+
return Promise.all(rows.map(async ({url, versions, datasets, updatedAt}) => {
112+
var datasetVersions = {}
113+
versions = versions.split(',')
114+
datasets = datasets.split(',')
115+
for (let i = 0; i < datasets.length; i++) {
116+
datasetVersions[datasets[i]] = Number(versions[i])
117+
}
118+
var meta = await archivesDb.getMeta(toHostname(url))
119+
return {url, title: meta.title, datasetVersions, updatedAt}
120+
}))
121+
}
122+
123+
const resetSite =
124+
exports.resetSite = async function (url) {
125+
await db.run(`DELETE FROM crawl_sources WHERE url = ?`, [url])
126+
}
127+
128+
exports.WEBAPI = {createEventsStream, getCrawlStates, resetSite}
129+
130+
// internal methods
131+
// =
132+
133+
function toHostname (url) {
134+
url = new URL(url)
135+
return url.hostname
136+
}

crawler/util.js

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,14 @@
1+
const EventEmitter = require('events')
12
const pump = require('pump')
23
const concat = require('concat-stream')
34
const db = require('../dbs/profile-data-db')
45
const dat = require('../dat')
56

67
const READ_TIMEOUT = 30e3
78

9+
const crawlerEvents = new EventEmitter()
10+
exports.crawlerEvents = crawlerEvents
11+
812
exports.doCrawl = async function (archive, crawlSource, crawlDataset, crawlDatasetVersion, handlerFn) {
913
const url = archive.url
1014

@@ -38,14 +42,19 @@ exports.doCrawl = async function (archive, crawlSource, crawlDataset, crawlDatas
3842
)
3943
})
4044

45+
crawlerEvents.emit('crawl-dataset-start', {sourceUrl: archive.url, crawlDataset, crawlRange: {start, end}})
46+
4147
// handle changes
4248
await handlerFn({changes, resetRequired})
4349

4450
// final checkpoint
4551
await doCheckpoint(crawlDataset, crawlDatasetVersion, crawlSource, version)
52+
53+
crawlerEvents.emit('crawl-dataset-finish', {sourceUrl: archive.url, crawlDataset, crawlRange: {start, end}})
4654
}
4755

4856
const doCheckpoint = exports.doCheckpoint = async function (crawlDataset, crawlDatasetVersion, crawlSource, crawlSourceVersion) {
57+
crawlerEvents.emit('crawl-dataset-progress', {sourceUrl: crawlSource.url, crawlDataset, crawledVersion: crawlSourceVersion})
4958
await db.run(`DELETE FROM crawl_sources_meta WHERE crawlDataset = ? AND crawlSourceId = ?`, [crawlDataset, crawlSource.id])
5059
await db.run(`
5160
INSERT

web-apis/bg.js

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ const downloadsManifest = require('./manifests/internal/downloads')
1010
const historyManifest = require('./manifests/internal/history')
1111
const sitedataManifest = require('./manifests/internal/sitedata')
1212
const watchlistManifest = require('./manifests/internal/watchlist')
13+
const crawlerManifest = require('./manifests/internal/crawler')
1314
const postsManifest = require('./manifests/internal/posts')
1415
const followgraphManifest = require('./manifests/internal/followgraph')
1516

@@ -19,6 +20,7 @@ const bookmarksAPI = require('./bg/bookmarks')
1920
const historyAPI = require('./bg/history')
2021
const sitedataAPI = require('../dbs/sitedata').WEBAPI
2122
const watchlistAPI = require('./bg/watchlist')
23+
const crawlerAPI = require('../crawler').WEBAPI
2224
const postsAPI = require('./bg/posts')
2325
const followgraphAPI = require('./bg/followgraph')
2426

@@ -54,6 +56,7 @@ exports.setup = function () {
5456
globals.rpcAPI.exportAPI('history', historyManifest, historyAPI, internalOnly)
5557
globals.rpcAPI.exportAPI('sitedata', sitedataManifest, sitedataAPI, internalOnly)
5658
globals.rpcAPI.exportAPI('watchlist', watchlistManifest, watchlistAPI, internalOnly)
59+
globals.rpcAPI.exportAPI('crawler', crawlerManifest, crawlerAPI, internalOnly)
5760
globals.rpcAPI.exportAPI('posts', postsManifest, postsAPI, internalOnly)
5861
globals.rpcAPI.exportAPI('followgraph', followgraphManifest, followgraphAPI, internalOnly)
5962

web-apis/fg/beaker.js

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ const downloadsManifest = require('../manifests/internal/downloads')
88
const historyManifest = require('../manifests/internal/history')
99
const sitedataManifest = require('../manifests/internal/sitedata')
1010
const watchlistManifest = require('../manifests/internal/watchlist')
11+
const crawlerManifest = require('../manifests/internal/crawler')
1112
const postsManifest = require('../manifests/internal/posts')
1213
const followgraphManifest = require('../manifests/internal/followgraph')
1314

@@ -24,6 +25,7 @@ exports.setup = function (rpc) {
2425
const historyRPC = rpc.importAPI('history', historyManifest, opts)
2526
const sitedataRPC = rpc.importAPI('sitedata', sitedataManifest, opts)
2627
const watchlistRPC = rpc.importAPI('watchlist', watchlistManifest, opts)
28+
const crawlerRPC = rpc.importAPI('crawler', crawlerManifest, opts)
2729
const postsRPC = rpc.importAPI('posts', postsManifest, opts)
2830
const followgraphRPC = rpc.importAPI('followgraph', followgraphManifest, opts)
2931

@@ -158,6 +160,12 @@ exports.setup = function (rpc) {
158160
beaker.watchlist.remove = watchlistRPC.remove
159161
beaker.watchlist.createEventsStream = () => fromEventStream(watchlistRPC.createEventsStream())
160162

163+
// beaker.crawler
164+
beaker.crawler = {}
165+
beaker.crawler.getCrawlStates = crawlerRPC.getCrawlStates
166+
beaker.crawler.resetSite = crawlerRPC.resetSite
167+
beaker.crawler.createEventsStream = () => fromEventStream(crawlerRPC.createEventsStream())
168+
161169
// beaker.posts
162170
beaker.posts = {}
163171
beaker.posts.list = postsRPC.list
Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
module.exports = {
2+
getCrawlStates: 'promise',
3+
resetSite: 'promise',
4+
createEventsStream: 'readable'
5+
}

0 commit comments

Comments
 (0)