Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
112 changes: 112 additions & 0 deletions backend/fetching/channels/ooni.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,112 @@
const { PollChannel } = require('downstream');
const Report = require('../../models/report');
const { fetchDailyMeasurements } = require('../ooniApi');
const { normalizeDailyCounts, evaluateAlert } = require('../ooniAlerts');

const DAY_MS = 24 * 60 * 60 * 1000;
const ALERT_DELAY_HOURS = 6;
const NETWORK_NAMES = {
44244: 'IranCell',
58224: 'MCCI',
};

function shiftDay(day, offset) {
return new Date(day.getTime() + offset * DAY_MS);
}

function dayString(day) {
return day.toISOString().slice(0, 10);
}

function alertDateFor(now) {
const date = new Date(now);
date.setUTCHours(0, 0, 0, 0);
if (now.getUTCHours() < ALERT_DELAY_HOURS) return shiftDay(date, -1);
return date;
}

function alertGuid(asn, alertDate) {
return `ooni:${asn}:volume:${alertDate}`;
}

function alertContent(asn, alerts) {
const network = NETWORK_NAMES[asn] || `AS${asn}`;
return `OONI volume alert for ${network} (AS${asn}): no web connectivity measurements were recorded on ${alerts[0].measurementDay}.`;
}

class OONIChannel extends PollChannel {
static INTERVAL = 60 * 60 * 1000;

constructor(options) {
const asns = String(options.asns || '')
.split(/[\s,]+/)
.filter(Boolean)
.map(Number);
if (asns.length === 0 || asns.some((asn) => !Number.isInteger(asn) || asn <= 0)) {
throw new Error('OONI sources require one or more valid ASNs.');
}

super({
...options,
namespace: options.namespace || `ooni-${asns.join('-')}`,
});
this.asns = asns;
this.interval = options.interval || OONIChannel.INTERVAL;
this.fetchDailyMeasurements = options.fetchDailyMeasurements || fetchDailyMeasurements;
}

async fetch() {
const alertDate = alertDateFor(new Date());
const since = dayString(shiftDay(alertDate, -1));
const until = dayString(alertDate);
const posts = [];

for (const asn of this.asns) {
const rows = await this.fetchDailyMeasurements({ asn, since, until });
const dailyCounts = normalizeDailyCounts(rows, since, until);
const alerts = evaluateAlert(dailyCounts, dayString(alertDate));
if (alerts.length === 0) continue;

const guid = alertGuid(asn, dayString(alertDate));
if (await Report.exists({ guid })) continue;

const post = this.parse({ asn, alerts, guid, fetchedAt: new Date() });
posts.push(post);
this.enqueue(post);
}

return posts;
}

parse(rawMessage) {
const { asn, alerts, guid, fetchedAt } = rawMessage;
const alertDate = alerts[0].alertDate;
const searchParams = new URLSearchParams({
probe_cc: 'IR',
probe_asn: `AS${asn}`,
test_name: 'web_connectivity',
since: alerts[0].measurementDay,
until: alertDate,
});

return {
authoredAt: new Date(`${alertDate}T00:00:00.000Z`),
fetchedAt,
author: `OONI AS${asn}`,
content: alertContent(asn, alerts),
url: `https://explorer.ooni.org/search?${searchParams}`,
platform: 'OONI',
platformID: guid,
raw: {
probeCC: 'IR',
probeASN: asn,
networkName: NETWORK_NAMES[asn] || null,
testName: 'web_connectivity',
alertDate,
triggers: alerts,
},
};
}
}

module.exports = OONIChannel;
2 changes: 1 addition & 1 deletion backend/fetching/hooks/postToReport.js
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@ module.exports = async function postToReport(post, next) {


let metadata;
if (platform === 'RSS') {
if (platform === 'RSS' || platform === 'OONI') {
// What do we want here?
metadata = {
// title: raw.title || null,
Expand Down
56 changes: 56 additions & 0 deletions backend/fetching/ooniAlerts.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
const DAY_MS = 24 * 60 * 60 * 1000;

function toDay(value) {
return new Date(`${value.slice(0, 10)}T00:00:00.000Z`);
}

function dayString(value) {
return value.toISOString().slice(0, 10);
}

function shiftDay(value, offset) {
return new Date(value.getTime() + offset * DAY_MS);
}

function normalizeDailyCounts(rows, since, until) {
const countsByDay = new Map();
rows.forEach((row) => {
const day = (row.measurement_start_day || '').slice(0, 10);
if (day) countsByDay.set(day, Number(row.measurement_count) || 0);
});

const normalized = [];
for (let day = toDay(since); day < toDay(until); day = shiftDay(day, 1)) {
const key = dayString(day);
normalized.push({
day: key,
measurementCount: countsByDay.get(key) || 0,
});
}
return normalized;
}

function evaluateAlert(dailyCounts, alertDate) {
const date = toDay(alertDate);
const countsByDay = new Map(
dailyCounts.map(({ day, measurementCount }) => [day, measurementCount]),
);
const latestDay = dayString(shiftDay(date, -1));
const latestCount = countsByDay.get(latestDay) || 0;

if (latestCount !== 0) return [];

return [
{
type: 'zero_measurements',
alertDate: dayString(date),
measurementDay: latestDay,
measurementCount: 0,
},
];
}

module.exports = {
normalizeDailyCounts,
evaluateAlert,
};
28 changes: 28 additions & 0 deletions backend/fetching/ooniApi.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
const AGGREGATION_URL = 'https://api.ooni.org/api/v1/aggregation';

async function fetchDailyMeasurements({ asn, since, until, fetchImpl = fetch }) {
const params = new URLSearchParams({
probe_cc: 'IR',
probe_asn: String(asn),
test_name: 'web_connectivity',
axis_x: 'measurement_start_day',
since,
until,
});
const url = `${AGGREGATION_URL}?${params}`;
const response = await fetchImpl(url, {
headers: { accept: 'application/json' },
});
if (!response.ok) {
throw new Error(`OONI aggregation request failed (${response.status}): ${url}`);
}

const payload = await response.json();
const result = payload.result || [];
return Array.isArray(result) ? result : [result];
}

module.exports = {
AGGREGATION_URL,
fetchDailyMeasurements,
};
8 changes: 8 additions & 0 deletions backend/fetching/sourceToChannel.js
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ const downstream = require('./downstream');
const AggieCrowdTangleChannel = require('./channels/crowdtangle');
const TelegramChannel = require('./channels/telegram');
const RSSChannel = require('./channels/rss');
const OONIChannel = require('./channels/ooni');

const { TwitterPageChannel, JunkipediaChannel } = builtin;

Expand Down Expand Up @@ -187,6 +188,13 @@ function createChannel(source) {
};
channel = new RSSChannel(options);
break;
case 'ooni':
options = {
...options,
asns: lists,
};
channel = new OONIChannel(options);
break;
default:
}

Expand Down
3 changes: 2 additions & 1 deletion backend/models/credentials.js
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,8 @@ const credentialsTypes = [
'crowdtangle',
'telegram',
'junkipedia',
'rss'
'rss',
'ooni'
];

// validates secrete based on their type
Expand Down
2 changes: 1 addition & 1 deletion backend/models/source.js
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ var urlValidator = function (url) {
)
}

var mediaValues = ['facebook', 'instagram', 'comments', 'elmo', 'twitter', 'rss', 'dummy', 'smsgh', 'whatsapp', 'telegram', 'junkipedia', 'dummy-pull', 'dummy-fast'];
var mediaValues = ['facebook', 'instagram', 'comments', 'elmo', 'twitter', 'rss', 'ooni', 'dummy', 'smsgh', 'whatsapp', 'telegram', 'junkipedia', 'dummy-pull', 'dummy-fast'];

var sourceSchema = new mongoose.Schema({
media: { type: String, enum: mediaValues },
Expand Down
44 changes: 44 additions & 0 deletions docs/OONI.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
# OONI integration

Aggie polls OONI's public aggregation API for Iran `web_connectivity`
measurement volume on AS44244 (IranCell) and AS58224 (MCCI). Alerts enter the
normal Aggie report pipeline and use deterministic daily GUIDs, so repeated
polling does not create duplicate reports.

## Alert rules

For an alert date `D`, only completed UTC days are evaluated:

Aggie creates a `zero_measurements` alert when the measurement count on `D-1`
is zero. OONI omits zero-count days from its response, so missing calendar dates
are filled with zero before evaluation. Production polling waits until 06:00 UTC
before evaluating the previous UTC day, reducing false alerts while OONI
finishes publishing daily aggregates.

## Configure Aggie

1. Open Settings, then Credentials.
2. Create an `ooni` credential. OONI is public, so no token is requested.
3. Open Settings, then Sources.
4. Create an `ooni` source and select the credential.
5. Keep the default ASN list `44244 58224`.
6. Enable the source and turn global fetching on.

OONI alerts appear as normal reports with media type `OONI`. Report metadata
contains the ASN, alert type, date, and zero measurement count. The report link
opens the corresponding OONI Explorer query.

## Historical backtest

Run the same evaluator used by the production channel:

```powershell
npm run backtest:ooni -- 2026-07-30
```

The optional date is the alert date through which to evaluate. Output is written
to the ignored `data/ooni-alert-backtest.json` and
`data/ooni-alert-backtest.csv` files.

The backtest emits only zero-measurement alerts. There is no additional cooldown
or incident-level suppression.
Loading